Project
Loading...
Searching...
No Matches
O2HitMerger.cxx
Go to the documentation of this file.
1// Copyright 2019-2020 CERN and copyright holders of ALICE O2.
2// See https://alice-o2.web.cern.ch/copyright for details of the copyright holders.
3// All rights not expressly granted are reserved.
4//
5// This software is distributed under the terms of the GNU General Public
6// License v3 (GPL Version 3), copied verbatim in the file "COPYING".
7//
8// In applying this license CERN does not waive the privileges and immunities
9// granted to it by virtue of its status as an Intergovernmental Organization
10// or submit itself to any jurisdiction.
11
13
14#include "O2HitMerger.h"
15#include <memory>
16#include <string>
17#include <type_traits>
18#include <fairmq/Message.h>
19#include <fairmq/Device.h>
20#include <fairlogger/Logger.h>
22#include <DetectorsBase/Stack.h>
26#include <gsl/gsl>
27#include <TFile.h>
28#include <TMemFile.h>
29#include <TTree.h>
30#include <TString.h>
31#include <TSystem.h>
32#include <TROOT.h>
33#include <memory>
34#include <TMessage.h>
35#include <fairmq/Parts.h>
36#include <ctime>
37#include <TStopwatch.h>
38#include <sstream>
39#include <cassert>
40#include "FairSystemInfo.h"
41
42#include "PrimaryServerState.h"
43#include <SimConfig/SimConfig.h>
48#include <TOFSimulation/Detector.h>
61
63#include <map>
64#include <vector>
65#include <list>
66#include <csignal>
67#include <mutex>
68#include <atomic>
69#include <filesystem>
70#include <functional>
71
73
74#ifdef ENABLE_UPGRADES
84#endif
85
86#include <tbb/concurrent_unordered_map.h>
87#include <thread>
88
89namespace o2
90{
91namespace devices
92{
93
94namespace
95{
96// Function communicating to primary particle server that it is now safe to shutdown.
97// From the perspective of o2-sim, this is the case when all configs have been propagated and the system
98// is running ok: For instance after the HitMerger is initialized and got it's first data from Geant workers.
99bool primaryServer_sendShutdownPermission(fair::mq::Channel& channel)
100{
101 std::unique_ptr<fair::mq::Message> request(channel.NewSimpleMessage((int)o2::O2PrimaryServerInfoRequest::AllowShutdown));
102 std::unique_ptr<fair::mq::Message> reply(channel.NewMessage());
103
104 int timeoutinMS = 100;
105 if (channel.Send(request, timeoutinMS) > 0) {
106 LOG(info) << "Sending Shutdown permission to particle server";
107 if (channel.Receive(reply, timeoutinMS) > 0) {
108 // the answer is a simple ack with a status code
109 LOG(info) << "Shutdown permission was acknowledged";
110 } else {
111 LOG(error) << "No answer received within " << timeoutinMS << "ms\n";
112 return false;
113 }
114 return true;
115 }
116 return false;
117}
118} // namespace
119
121{
122 mTimer.Start();
123 mInitialOutputDir = std::filesystem::current_path().string();
124 mCurrentOutputDir = mInitialOutputDir;
125}
126
128{
129 FairSystemInfo sysinfo;
130 LOG(info) << "TIME-STAMP " << mTimer.RealTime() << "\t";
131 mTimer.Continue();
132 LOG(info) << "MEM-STAMP " << sysinfo.GetCurrentMemory() / (1024. * 1024) << " "
133 << sysinfo.GetMaxMemory() << " MB\n";
134}
135
136void O2HitMerger::InitTask()
137{
138 LOG(info) << "INIT HIT MERGER";
139 ROOT::EnableThreadSafety();
140
141 std::string outfilename("o2sim_merged_hits.root"); // default name
142 // query the sim config ... which is used to extract the filenames
143 if (o2::querySimConfig(GetChannels().at("o2sim-primserv-info").at(0))) {
145 mNExpectedEvents = o2::conf::SimConfig::Instance().getNEvents();
146 } else {
147 // we didn't manage to get a configuration --> better to fail
148 LOG(fatal) << "No configuration received. Aborting";
149 }
153
154 mOutFileName = outfilename.c_str();
155 if (mWriteToDisc) {
156 mOutFile = new TFile(outfilename.c_str(), "RECREATE");
157 mOutTree = new TTree("o2sim", "o2sim");
158 mOutTree->SetDirectory(mOutFile);
159
160 mMCHeaderOnlyOutFile = new TFile(o2::base::NameConf::getMCHeadersFileName(o2::conf::SimConfig::Instance().getOutPrefix().c_str()).c_str(), "RECREATE");
161 mMCHeaderTree = new TTree("o2sim", "o2sim");
162 mMCHeaderTree->SetDirectory(mMCHeaderOnlyOutFile);
163 }
164 // detectors init only once
165 if (mDetectorInstances.size() == 0) {
166 initDetInstances();
167 // has to be after init of Detectors
169 initHitFiles(o2::conf::SimConfig::Instance().getOutPrefix());
170 }
171
172 // init pipe
173 auto pipeenv = getenv("ALICE_O2SIMMERGERTODRIVER_PIPE");
174 if (pipeenv) {
175 mPipeToDriver = atoi(pipeenv);
176 LOG(info) << "ASSIGNED PIPE HANDLE " << mPipeToDriver;
177 } else {
178 LOG(warning) << "DID NOT FIND ENVIRONMENT VARIABLE TO INIT PIPE";
179 }
180
181 // if no data to expect we shut down the device NOW since it would otherwise hang
182 if (mNExpectedEvents == 0) {
183 if (mAsService) {
184 waitForControlInput();
185 } else {
186 LOG(info) << "NOT EXPECTING ANY DATA; SHUTTING DOWN";
187 raise(SIGINT);
188 }
189 }
190}
191
192bool O2HitMerger::setWorkingDirectory(std::string const& dir)
193{
194 namespace fs = std::filesystem;
195
196 // sets the output directory where simulation files are produced
197 // and creates it when it doesn't exist already
198
199 // 2 possibilities:
200 // a) dir is relative dir. Then we interpret it as relative to the initial
201 // base directory
202 // b) or dir is itself absolut.
203 try {
204 fs::current_path(fs::path(mInitialOutputDir)); // <--- to make sure relative start is always the same
205 if (!dir.empty()) {
206 auto absolutePath = fs::absolute(fs::path(dir));
207 if (!fs::exists(absolutePath)) {
208 if (!fs::create_directory(absolutePath)) {
209 LOG(error) << "Could not create directory " << absolutePath.string();
210 return false;
211 }
212 }
213 // set the current path
214 fs::current_path(absolutePath.string().c_str());
215 mCurrentOutputDir = fs::current_path().string();
216 }
217 LOG(info) << "FINAL PATH " << mCurrentOutputDir;
218 } catch (std::exception e) {
219 LOG(error) << " could not change path to " << dir;
220 }
221 return true;
222}
223
224bool O2HitMerger::ReInit(o2::conf::SimReconfigData const& reconfig)
225{
226 if (reconfig.stop) {
227 return false;
228 }
229 if (!setWorkingDirectory(reconfig.outputDir)) {
230 return false;
231 }
232
233 std::string outfilename("o2sim_merged_hits.root"); // default name
235 mNExpectedEvents = reconfig.nEvents;
236 mOutFileName = outfilename.c_str();
237 if (mWriteToDisc) {
238 mOutFile = new TFile(outfilename.c_str(), "RECREATE");
239 mOutTree = new TTree("o2sim", "o2sim");
240 mOutTree->SetDirectory(mOutFile);
241
242 mMCHeaderOnlyOutFile = new TFile(o2::base::NameConf::getMCHeadersFileName(reconfig.outputPrefix).c_str(), "RECREATE");
243 mMCHeaderTree = new TTree("o2sim", "o2sim");
244 mMCHeaderTree->SetDirectory(mMCHeaderOnlyOutFile);
245 }
246 // reinit detectorInstance files (also make sure they are closed before continuing)
247 initHitFiles(reconfig.outputPrefix);
248
249 // clear "counter" datastructures
250 mPartsCheckSum.clear();
251 mEventChecksum = 0;
252
253 // clear collector datastructures
254 mMCTrackBuffer.clear();
255 mTrackRefBuffer.clear();
256 mSubEventInfoBuffer.clear();
257 mFlushableEvents.clear();
258 mNextFlushID = 1;
259
260 return true;
261}
262
263template <typename T, typename V>
264V O2HitMerger::insertAdd(std::map<T, V>& m, T const& key, V value)
265{
266 const auto iter = m.find(key);
267 V accum{0};
268 if (iter != m.end()) {
269 iter->second += value;
270 accum = iter->second;
271 } else {
272 m.insert(std::make_pair(key, value));
273 accum = value;
274 }
275 return accum;
276}
277
278template <typename T>
279bool O2HitMerger::isDataComplete(T checksum, T nparts)
280{
281 return checksum == nparts * (nparts + 1) / 2;
282}
283
284void O2HitMerger::consumeHits(int eventID, fair::mq::Parts& data, int& index)
285{
286 auto headermessage = std::move(data.At(index++));
287 // this should be the header announcing the hits of one detector
288 if (headermessage->GetSize() == sizeof(o2::base::HitsHeader)) {
289 auto header = *static_cast<o2::base::HitsHeader const*>(headermessage->GetData());
291 LOG(debug2) << "I1 " << header.detID << " NAME " << id.getName() << " MB "
292 << data.At(index)->GetSize() / 1024. / 1024.;
293
294 // get the detector that can interpret it
295 auto detector = mDetectorInstances[id].get();
296 if (detector) {
297 detector->collectHits(eventID, data, index, header.shm);
298 }
299 }
300}
301
302template <typename T, typename BT>
303void O2HitMerger::consumeData(int eventID, fair::mq::Parts& data, int& index, BT& buffer)
304{
305 auto decodeddata = o2::base::decodeTMessage<T*>(data, index);
306 if (buffer.find(eventID) == buffer.end()) {
307 buffer[eventID] = typename BT::mapped_type();
308 }
309 buffer[eventID].push_back(decodeddata);
310 // delete decodeddata; --> we store the pointers
311 index++;
312}
313
314void O2HitMerger::fillSubEventInfoEntry(o2::data::SubEventInfo& info)
315{
316 if (mSubEventInfoBuffer.find(info.eventID) == mSubEventInfoBuffer.end()) {
317 mSubEventInfoBuffer[info.eventID] = std::list<o2::data::SubEventInfo*>();
318 }
319 mSubEventInfoBuffer[info.eventID].push_back(&info);
320}
321
322bool O2HitMerger::waitForControlInput()
323{
324 o2::simpubsub::publishMessage(GetChannels()["merger-notifications"].at(0), o2::simpubsub::simStatusString("MERGER", "STATUS", "AWAITING INPUT"));
325
326 auto factory = fair::mq::TransportFactory::CreateTransportFactory("zeromq");
327 auto channel = fair::mq::Channel{"o2sim-control", "sub", factory};
328 auto controlsocketname = getenv("ALICE_O2SIMCONTROL");
329 LOG(info) << "SOCKETNAME " << controlsocketname;
330 channel.Connect(std::string(controlsocketname));
331 channel.Validate();
332 std::unique_ptr<fair::mq::Message> reply(channel.NewMessage());
333
334 LOG(info) << "WAITING FOR INPUT";
335 if (channel.Receive(reply) > 0) {
336 auto data = reply->GetData();
337 auto size = reply->GetSize();
338
339 std::string command(reinterpret_cast<char const*>(data), size);
340 LOG(info) << "message: " << command;
341
343 o2::conf::parseSimReconfigFromString(command, reconfig);
344 return ReInit(reconfig);
345 } else {
346 LOG(info) << "NOTHING RECEIVED";
347 }
348 return true;
349}
350
351bool O2HitMerger::ConditionalRun()
352{
353 auto& channel = GetChannels().at("simdata").at(0);
354 fair::mq::Parts request;
355 auto bytes = channel.Receive(request);
356 if (bytes < 0) {
357 LOG(error) << "Some error occurred on socket during receive on sim data";
358 return true; // keep going
359 }
360 TStopwatch timer;
361 timer.Start();
362 auto more = handleSimData(request, 0);
363 LOG(info) << "HitMerger processing took " << timer.RealTime();
364 if (!more && mAsService) {
365 LOG(info) << " CONTROL ";
366 // if we are done treating data we may go back to init phase
367 // for the next batch
368 return waitForControlInput();
369 }
370
371 static bool initAcknowledged = false;
372 if (!initAcknowledged) {
373 primaryServer_sendShutdownPermission(GetChannels().at("o2sim-primserv-info").at(0));
374 initAcknowledged = true;
375 }
376
377 return more;
378}
379
380bool O2HitMerger::handleSimData(fair::mq::Parts& data, int /*index*/)
381{
382 bool expectmore = true;
383 int index = 0;
384 auto infoptr = o2::base::decodeTMessage<o2::data::SubEventInfo*>(data, index++);
385 o2::data::SubEventInfo& info = *infoptr;
386 // once a merge thread runs, the buffered info of a complete event may be freed at any time
387 const auto eventID = info.eventID;
388 const auto maxEvents = info.maxEvents;
389 auto accum = insertAdd<uint32_t, uint32_t>(mPartsCheckSum, info.eventID, (uint32_t)info.part);
390
391 LOG(info) << "SIMDATA channel got " << data.Size() << " parts for event " << info.eventID << " part " << info.part << " out of " << info.nparts;
392
393 fillSubEventInfoEntry(info);
394 consumeData<std::vector<o2::MCTrack>>(info.eventID, data, index, mMCTrackBuffer);
395 consumeData<std::vector<o2::TrackReference>>(info.eventID, data, index, mTrackRefBuffer);
396 while (index < data.Size()) {
397 consumeHits(info.eventID, data, index);
398 }
399
400 if (isDataComplete<uint32_t>(accum, info.nparts)) {
401 LOG(info) << "Event " << info.eventID << " complete. Marking as flushable";
402 mFlushableEvents[info.eventID] = true;
403
404 // check if previous flush finished
405 // start merging only when no merging currently happening
406 // Like this we don't have to join/wait on the thread here and do not block the outer ConditionalRun handling
407 // TODO: Let this run fully asynchronously (not even triggered by ConditionalRun)
408 if (!mergingInProgress) {
409 if (mMergerIOThread.joinable()) {
410 mMergerIOThread.join();
411 }
412 // start hit merging and flushing in a separate thread in order not to block
413 mMergerIOThread = std::thread([this]() { mergingInProgress = true; mergeAndFlushData(); mergingInProgress = false; });
414 }
415
416 mEventChecksum += eventID;
417 // we also need to check if we have all events
418 if (isDataComplete<uint32_t>(mEventChecksum, maxEvents)) {
419 LOG(info) << "ALL EVENTS HERE; CHECKSUM " << mEventChecksum;
420
421 // flush remaining data and close file
422 if (mMergerIOThread.joinable()) {
423 mMergerIOThread.join();
424 }
425 mMergerIOThread = std::thread([this]() { mergingInProgress = true; mergeAndFlushData(); mergingInProgress = false; });
426 if (mMergerIOThread.joinable()) {
427 mMergerIOThread.join();
428 }
429
430 expectmore = false;
431 }
432
433 if (mPipeToDriver != -1) {
434 if (write(mPipeToDriver, &eventID, sizeof(eventID)) == -1) {
435 LOG(error) << "FAILED WRITING TO PIPE";
436 };
437 }
438 }
439 return expectmore;
440}
441
442void O2HitMerger::cleanEvent(int eventID)
443{
444 auto release = [eventID](auto& buffer) {
445 auto iter = buffer.find(eventID);
446 if (iter != buffer.end()) {
447 for (auto ptr : iter->second) {
448 delete ptr;
449 }
450 iter->second = {};
451 }
452 };
453 release(mMCTrackBuffer);
454 release(mTrackRefBuffer);
455 release(mSubEventInfoBuffer);
456}
457
458template <typename T>
459void O2HitMerger::backInsert(T const& from, T& to)
460{
461 std::copy(from.begin(), from.end(), std::back_inserter(to));
462}
463
464void O2HitMerger::reorderAndMergeMCTracks(int eventID, TTree* target, const std::vector<int>& nprimaries, const std::vector<int>& nsubevents, std::function<void(std::vector<MCTrack> const&)> tracks_analysis_hook, o2::dataformats::MCEventHeader const* mceventheader)
465{
466 // avoid doing this for trivial cases
467 std::vector<MCTrack>* mcTracksPerSubEvent = nullptr;
468 auto targetdata = std::make_unique<std::vector<MCTrack>>();
469
470 auto& vectorOfSubEventMCTracks = mMCTrackBuffer[eventID];
471 const auto entries = vectorOfSubEventMCTracks.size();
472
473 if (entries > 1) {
474 size_t ntracks = 0;
475 for (auto tracks : vectorOfSubEventMCTracks) {
476 ntracks += tracks->size();
477 }
478 targetdata->reserve(ntracks);
479 //
480 // loop over subevents to store the primary events
481 //
482 int nprimTot = 0;
483 for (int entry = entries - 1; entry >= 0; --entry) {
484 int index = nsubevents[entry];
485 nprimTot += nprimaries[index];
486 for (int i = 0; i < nprimaries[index]; i++) {
487 auto& track = (*vectorOfSubEventMCTracks[index])[i];
488 if (track.isTransported()) { // reset daughters only if track was transported, it will be fixed below
489 track.SetFirstDaughterTrackId(-1);
490 track.SetLastDaughterTrackId(-1);
491 }
492 targetdata->push_back(track);
493 }
494 }
495 //
496 // loop a second time to store the secondaries and fix the mother track IDs
497 //
498 Int_t idelta1 = nprimTot;
499 Int_t idelta0 = 0;
500 for (int entry = entries - 1; entry >= 0; --entry) {
501 int index = nsubevents[entry];
502
503 auto& subEventTracks = *(vectorOfSubEventMCTracks[index]);
504 // we need to fetch the right mctracks here!!
505 Int_t npart = (int)(subEventTracks.size());
506 Int_t nprim = nprimaries[index];
507 idelta1 -= nprim;
508
509 for (Int_t i = nprim; i < npart; i++) {
510 auto& track = subEventTracks[i];
511 Int_t cId = track.getMotherTrackId();
512 if (cId >= nprim) {
513 cId += idelta1;
514 } else {
515 cId += idelta0;
516 }
517 track.SetMotherTrackId(cId);
518 track.SetFirstDaughterTrackId(-1);
519
520 Int_t hwm = (int)(targetdata->size());
521 auto& mother = (*targetdata)[cId];
522 if (mother.getFirstDaughterTrackId() == -1) {
523 mother.SetFirstDaughterTrackId(hwm);
524 }
525 mother.SetLastDaughterTrackId(hwm);
526
527 targetdata->push_back(track);
528 }
529 idelta0 += nprim;
530 idelta1 += npart;
531 }
532 }
533 //
534 // write to output
535 auto filladdr = (entries > 1) ? targetdata.get() : vectorOfSubEventMCTracks[0];
536
537 // we give the possibility to produce some MC track statistics
538 // to be saved as part of the MCHeader structure
539 tracks_analysis_hook(*filladdr);
540
541 if (mWriteToDisc && target) {
542 auto targetbr = o2::base::getOrMakeBranch(*target, "MCTrack", &filladdr);
543 targetbr->SetAddress(&filladdr);
544 targetbr->Fill();
545 targetbr->ResetAddress();
546 }
547 // forwarding the track data to other consumers (pub/sub)
548 if (mForwardKine) {
549 auto free_tmessage = [](void* data, void* hint) { delete static_cast<TMessage*>(hint); };
550 auto& channel = GetChannels().at("kineforward").at(0);
551 TMessage* tmsg = new TMessage(kMESS_OBJECT);
552 tmsg->WriteObjectAny((void*)filladdr, TClass::GetClass("std::vector<o2::MCTrack>"));
553 std::unique_ptr<fair::mq::Message> trackmessage(channel.NewMessage(tmsg->Buffer(), tmsg->BufferSize(), free_tmessage, tmsg));
554 tmsg = new TMessage(kMESS_OBJECT);
555 tmsg->WriteObjectAny((void*)mceventheader, TClass::GetClass("o2::dataformats::MCEventHeader"));
556 std::unique_ptr<fair::mq::Message> headermessage(channel.NewMessage(tmsg->Buffer(), tmsg->BufferSize(), free_tmessage, tmsg));
557 fair::mq::Parts reply;
558 reply.AddPart(std::move(headermessage));
559 reply.AddPart(std::move(trackmessage));
560 channel.Send(reply);
561 LOG(info) << "Forward publish MC tracks on channel";
562 }
563}
564
565template <typename T, typename M>
566void O2HitMerger::remapTrackIdsAndMerge(std::string brname, int eventID, TTree& target, const std::vector<int>& trackoffsets, const std::vector<int>& nprimaries, const std::vector<int>& subevOrdered, M& mapOfVectorOfTs)
567{
568 //
569 // Remap the mother track IDs by adding an offset.
570 // The offset calculated as the sum of the number of entries in the particle list of the previous subevents.
571 // This method is called by O2HitMerger::mergeAndFlushData(int)
572 //
573 T* incomingdata = nullptr;
574 std::unique_ptr<T> targetdata(nullptr);
575 auto& vectorOfT = mapOfVectorOfTs[eventID];
576 const auto entries = vectorOfT.size();
577
578 if (entries == 1) {
579 // nothing to do in case there is only one entry
580 incomingdata = vectorOfT[0];
581 } else {
582 targetdata = std::make_unique<T>();
583 size_t nentries = 0;
584 for (auto data : vectorOfT) {
585 nentries += data->size();
586 }
587 targetdata->reserve(nentries);
588 // loop over subevents
589 Int_t nprimTot = 0;
590 for (int entry = 0; entry < entries; entry++) {
591 nprimTot += nprimaries[entry];
592 }
593 Int_t idelta0 = 0;
594 Int_t idelta1 = nprimTot;
595 for (int entry = entries - 1; entry >= 0; --entry) {
596 Int_t index = subevOrdered[entry];
597 Int_t nprim = nprimaries[index];
598 incomingdata = vectorOfT[index];
599 idelta1 -= nprim;
600 for (auto& data : *incomingdata) {
601 updateTrackIdWithOffset(data, nprim, idelta0, idelta1);
602 targetdata->push_back(data);
603 }
604 idelta0 += nprim;
605 idelta1 += trackoffsets[index];
606 }
607 }
608 auto dataaddr = (entries == 1) ? incomingdata : targetdata.get();
609 auto targetbr = o2::base::getOrMakeBranch(target, brname.c_str(), &dataaddr);
610 targetbr->SetAddress(&dataaddr);
611 targetbr->Fill();
612 targetbr->ResetAddress();
613}
614
615void O2HitMerger::updateTrackIdWithOffset(MCTrack& track, Int_t nprim, Int_t idelta0, Int_t idelta1)
616{
617 Int_t cId = track.getMotherTrackId();
618 Int_t ioffset = (cId < nprim) ? idelta0 : idelta1;
619 if (cId != -1) {
620 track.SetMotherTrackId(cId + ioffset);
621 }
622}
623
624void O2HitMerger::updateTrackIdWithOffset(TrackReference& ref, Int_t nprim, Int_t idelta0, Int_t idelta1)
625{
626 ref.setTrackID(o2::base::Detector::offsetTrackIndex(ref.getTrackID(), nprim, idelta0, idelta1));
627}
628
629void O2HitMerger::initHitTreeAndOutFile(std::string prefix, int detID)
630{
632 if (mDetectorOutFiles.find(detID) != mDetectorOutFiles.end() && mDetectorOutFiles[detID]) {
633 LOG(warn) << "Hit outfile for detID " << DetID::getName(detID) << " already initialized --> Reopening";
634 mDetectorOutFiles[detID]->Close();
635 delete mDetectorOutFiles[detID];
636 }
637 std::string name(o2::base::DetectorNameConf::getHitsFileName(detID, prefix));
638 if (mWriteToDisc) {
639 mDetectorOutFiles[detID] = new TFile(name.c_str(), "RECREATE");
640 mDetectorToTTreeMap[detID] = new TTree("o2sim", "o2sim");
641 mDetectorToTTreeMap[detID]->SetDirectory(mDetectorOutFiles[detID]);
642 } else {
643 mDetectorOutFiles[detID] = nullptr;
644 mDetectorToTTreeMap[detID] = nullptr;
645 }
646}
647
648bool O2HitMerger::mergeAndFlushData()
649{
650 auto isFlushable = [this](int eventID) {
651 auto iter = mFlushableEvents.find(eventID);
652 return iter != mFlushableEvents.end() && iter->second;
653 };
654
655 LOG(info) << "Launching merge kernel ";
656 if (!isFlushable(mNextFlushID)) {
657 return false;
658 }
659 for (; isFlushable(mNextFlushID); ++mNextFlushID) {
660 auto flusheventID = mNextFlushID;
661 LOG(info) << "Merge and flush event " << flusheventID;
662 auto iter = mSubEventInfoBuffer.find(flusheventID);
663 if (iter == mSubEventInfoBuffer.end() || iter->second.size() == 0 || mNExpectedEvents == 0) {
664 LOG(error) << "No data entries found for event " << flusheventID;
665 continue;
666 }
667 auto& subEventInfoList = iter->second;
668
669 TStopwatch timer;
670 timer.Start();
671
672 // calculate trackoffsets
673 auto& confref = o2::conf::SimConfig::Instance();
674
675 // collecting trackoffsets (per data arrival id) to be used for global track-ID correction pass
676 std::vector<int> trackoffsets;
677 // collecting primary particles in each subevent (data arrival id)
678 std::vector<int> nprimaries;
679 // mapping of id to actual sub-event id (or part)
680 std::vector<int> nsubevents;
681
682 o2::dataformats::MCEventHeader* eventheader = nullptr; // The event header
683
684 // the MC labels (trackID) for hits
685 for (auto info : subEventInfoList) {
686 assert(info->npersistenttracks >= 0);
687 trackoffsets.emplace_back(info->npersistenttracks);
688 nprimaries.emplace_back(info->nprimarytracks);
689 nsubevents.emplace_back(info->part);
690 if (eventheader == nullptr) {
691 eventheader = &info->mMCEventHeader;
692 } else {
693 eventheader->getMCEventStats().add(info->mMCEventHeader.getMCEventStats());
694 }
695 }
696
697 // now see which events can be discarded in any case due to no hits
698 if (confref.isFilterOutNoHitEvents()) {
699 if (eventheader && eventheader->getMCEventStats().getNHits() == 0) {
700 LOG(info) << " Taking out event " << flusheventID << " due to no hits ";
701 cleanEvent(flusheventID);
702 continue;
703 }
704 }
705
706 // attention: We need to make sure that we write everything in the same event order
707 // but iteration over keys of a standard map in C++ is ordered
708
709 // b) merge the general data
710 //
711 // for MCTrack remap the motherIds and merge at the same go
712 const auto entries = subEventInfoList.size();
713 std::vector<int> subevOrdered((int)(nsubevents.size()));
714 for (int entry = entries - 1; entry >= 0; --entry) {
715 subevOrdered[nsubevents[entry] - 1] = entry;
716 }
717
718 // This is a hook that collects some useful statistics/properties on the event
719 // for use by other components;
720 // Properties are attached making use of the extensible "Info" feature which is already
721 // part of MCEventHeader. In such a way, one can also do this pass outside and attach arbitrary
722 // metadata to MCEventHeader without needing to change the data layout or API of the class itself.
723 // NOTE: This function might also be called directly in the primary server!?
724 auto mcheaderhook = [eventheader](std::vector<MCTrack> const& tracks) {
725 int eta1Point2Counter = 0;
726 int eta1Point0Counter = 0;
727 int eta0Point8Counter = 0;
728 int eta1Point2CounterPi = 0;
729 int eta1Point0CounterPi = 0;
730 int eta0Point8CounterPi = 0;
731 int prims = 0;
732 for (auto& tr : tracks) {
733 if (tr.isPrimary()) {
734 prims++;
735 const auto eta = tr.GetEta();
736 if (eta < 1.2) {
737 eta1Point2Counter++;
738 if (std::abs(tr.GetPdgCode()) == 211) {
739 eta1Point2CounterPi++;
740 }
741 }
742 if (eta < 1.0) {
743 eta1Point0Counter++;
744 if (std::abs(tr.GetPdgCode()) == 211) {
745 eta1Point0CounterPi++;
746 }
747 }
748 if (eta < 0.8) {
749 eta0Point8Counter++;
750 if (std::abs(tr.GetPdgCode()) == 211) {
751 eta0Point8CounterPi++;
752 }
753 }
754 } else {
755 break; // track layout is such that all prims are first anyway
756 }
757 }
758 // attach these properties to eventheader
759 // we only need to make the names standard
760 eventheader->putInfo("prims_eta_1.2", eta1Point2Counter);
761 eventheader->putInfo("prims_eta_1.0", eta1Point0Counter);
762 eventheader->putInfo("prims_eta_0.8", eta0Point8Counter);
763 eventheader->putInfo("prims_eta_1.2_pi", eta1Point2CounterPi);
764 eventheader->putInfo("prims_eta_1.0_pi", eta1Point0CounterPi);
765 eventheader->putInfo("prims_eta_0.8_pi", eta0Point8CounterPi);
766 eventheader->putInfo("prims_total", prims);
767 };
768 // the kinematics and each detector go to separate files, so we merge and flush them concurrently
769 // plain threads, since the TBB pool teardown crashed the merger
770 std::vector<std::thread> tasks;
771 tasks.emplace_back([&]() {
772 reorderAndMergeMCTracks(flusheventID, mOutTree, nprimaries, subevOrdered, mcheaderhook, eventheader);
773
774 if (mOutTree) {
775 // adjusting and merging track references
776 remapTrackIdsAndMerge<std::vector<o2::TrackReference>>("TrackRefs", flusheventID, *mOutTree, trackoffsets, nprimaries, subevOrdered, mTrackRefBuffer);
777
778 // write MC event headers
779 for (auto tree : {mOutTree, mMCHeaderTree}) {
780 auto headerbr = o2::base::getOrMakeBranch(*tree, "MCEventHeader.", &eventheader);
781 headerbr->SetAddress(&eventheader);
782 headerbr->Fill();
783 headerbr->ResetAddress();
784 }
785
786 // increase the entry count in the trees
787 mOutTree->SetEntries(mOutTree->GetEntries() + 1);
788 mMCHeaderTree->SetEntries(mMCHeaderTree->GetEntries() + 1);
789 }
790 });
791
792 // c) do the merge procedure for all hits ... delegate this to detector specific functions
793 // since they know about types; number of branches; etc.
794 // this will also fix the trackIDs inside the hits
795 for (int id = 0; id < mDetectorInstances.size(); ++id) {
796 auto& det = mDetectorInstances[id];
797 auto hittree = det ? mDetectorToTTreeMap[id] : nullptr;
798 if (hittree) {
799 tasks.emplace_back([&, det = det.get(), hittree]() {
800 det->mergeHitEntriesAndFlush(flusheventID, *hittree, trackoffsets, nprimaries, subevOrdered);
801 hittree->SetEntries(hittree->GetEntries() + 1);
802 });
803 }
804 }
805 for (auto& t : tasks) {
806 t.join();
807 }
808
809 cleanEvent(flusheventID);
810 LOG(info) << "Merge/flush for event " << flusheventID << " took " << timer.RealTime();
811 }
812 if (mWriteToDisc && mOutFile) {
813 LOG(info) << "Writing TTrees";
814 std::vector<TFile*> files{mOutFile, mMCHeaderOnlyOutFile};
815 for (int id = 0; id < mDetectorInstances.size(); ++id) {
816 if (mDetectorInstances[id] && mDetectorOutFiles[id]) {
817 files.push_back(mDetectorOutFiles[id]);
818 }
819 }
820 std::vector<std::thread> writers;
821 for (auto file : files) {
822 writers.emplace_back([file]() { file->Write("", TObject::kOverwrite); });
823 }
824 for (auto& t : writers) {
825 t.join();
826 }
827 }
828 return true;
829}
830
831void O2HitMerger::initHitFiles(std::string prefix)
832{
834
835 // a little helper lambda
836 auto isActivated = [](std::string s) -> bool {
837 // access user configuration for list of wanted modules
839 auto active = std::find(modulelist.begin(), modulelist.end(), s) != modulelist.end();
840 return active; };
841
842 for (int i = DetID::First; i <= DetID::Last; ++i) {
843 if (!isActivated(DetID::getName(i))) {
844 continue;
845 }
846 // init the detector specific output files
847 initHitTreeAndOutFile(prefix, i);
848 }
849
850 // external (CAD) detectors are not part of the readout-detector list (their module names
851 // are not DetID names); their slots were determined in initDetInstances()
852 for (auto detID : mExternalDetIDs) {
853 initHitTreeAndOutFile(prefix, detID);
854 }
855}
856
857// init detector instances used to write hit data to a TTree
858void O2HitMerger::initDetInstances()
859{
861
862 // a little helper lambda
863 auto isActivated = [](std::string s) -> bool {
864 // access user configuration for list of wanted modules
866 auto active = std::find(modulelist.begin(), modulelist.end(), s) != modulelist.end();
867 return active; };
868
869 // readout-only detector instances able to interpret and write the hits of each detector
870 using Factory = std::function<std::unique_ptr<o2::base::Detector>()>;
871 const std::map<int, Factory> factories{
872 {DetID::TPC, [] { return std::make_unique<o2::tpc::Detector>(true); }},
873 {DetID::ITS, [] { return std::make_unique<o2::its::Detector>(true); }},
874 {DetID::MFT, [] { return std::make_unique<o2::mft::Detector>(true); }},
875 {DetID::TRD, [] { return std::make_unique<o2::trd::Detector>(true); }},
876 {DetID::PHS, [] { return std::make_unique<o2::phos::Detector>(true); }},
877 {DetID::CPV, [] { return std::make_unique<o2::cpv::Detector>(true); }},
878 {DetID::EMC, [] { return std::make_unique<o2::emcal::Detector>(true); }},
879 {DetID::HMP, [] { return std::make_unique<o2::hmpid::Detector>(true); }},
880 {DetID::TOF, [] { return std::make_unique<o2::tof::Detector>(true); }},
881 {DetID::FT0, [] { return std::make_unique<o2::ft0::Detector>(true); }},
882 {DetID::FV0, [] { return std::make_unique<o2::fv0::Detector>(true); }},
883 {DetID::FDD, [] { return std::make_unique<o2::fdd::Detector>(true); }},
884 {DetID::MCH, [] { return std::make_unique<o2::mch::Detector>(true); }},
885 {DetID::MID, [] { return std::make_unique<o2::mid::Detector>(true); }},
886 {DetID::ZDC, [] { return std::make_unique<o2::zdc::Detector>(true); }},
887 {DetID::FOC, [] {
888 TString sName = "$O2_ROOT/share/Detectors/Geometry/FOC/geometryFiles/geometry_Sheets.txt";
889 gSystem->ExpandPathName(sName);
890 return std::make_unique<o2::focal::Detector>(true, sName.Data());
891 }},
892#ifdef ENABLE_UPGRADES
893 {DetID::IT3, [] { return std::make_unique<o2::its::Detector>(true, "IT3"); }},
894 {DetID::TRK, [] { return std::make_unique<o2::trk::Detector>(true); }},
895 {DetID::FT3, [] { return std::make_unique<o2::ft3::Detector>(true); }},
896 {DetID::FCT, [] { return std::make_unique<o2::fct::Detector>(true); }},
897 {DetID::TF3, [] { return std::make_unique<o2::iotof::Detector>(true); }},
898 {DetID::RCH, [] { return std::make_unique<o2::rich::Detector>(true); }},
899 {DetID::MI3, [] { return std::make_unique<o2::mi3::Detector>(true); }},
900 {DetID::ECL, [] { return std::make_unique<o2::ecal::Detector>(true); }},
901 {DetID::FD3, [] { return std::make_unique<o2::fd3::Detector>(true); }},
902#endif
903 };
904
905 mDetectorInstances.resize(DetID::nDetectors);
906 for (int i = DetID::First; i <= DetID::Last; ++i) {
907 if (!isActivated(DetID::getName(i))) {
908 continue;
909 }
910 auto factory = factories.find(i);
911 if (factory == factories.end()) {
912 LOG(warning) << "O2HitMerger: no hit merging available for readout detector " << DetID::getName(i);
913 continue;
914 }
915 mDetectorInstances[i] = factory->second();
916 }
917
918 // also register external (CAD-derived) sensitive detectors so their hits are persisted
919 // in parallel (multi-worker) mode
920 initExternalDetInstances();
921}
922
923// init detector instances for external (CAD-derived) sensitive detectors.
924// These are not part of the hard-coded DetID switch above: they are described in the
925// external geometry JSON (the same file used by build_geometry.C on the worker side) and
926// tied to an existing (free) DetID. The merger only needs an instance able to interpret the
927// generic o2::ext::Hit wire format and write the "<name>Hit" branch; no geometry is built here.
928void O2HitMerger::initExternalDetInstances()
929{
931
932 auto& simConfig = o2::conf::SimConfig::Instance();
933 const auto extGeomFile = simConfig.getExtGeomFilename();
934 if (extGeomFile.empty()) {
935 return;
936 }
937
938 // mirror the worker-side activation: an external detector participates when its module
939 // name is part of the active module list
940 auto const& activeModules = simConfig.getActiveModules();
941 auto isActivated = [&activeModules](std::string const& s) -> bool {
942 return std::find(activeModules.begin(), activeModules.end(), s) != activeModules.end();
943 };
944
945 for (auto* extdet : o2::ext::ExternalDetector::createFromJSON(extGeomFile)) {
946 const std::string name = extdet->GetName();
947 if (!isActivated(name)) {
948 delete extdet; // not requested in the active module list
949 continue;
950 }
951 const int detID = extdet->GetDetId();
952 if (detID < DetID::First || detID > DetID::Last) {
953 LOG(error) << "O2HitMerger: external detector " << name << " has invalid DetID " << detID << "; skipping";
954 delete extdet;
955 continue;
956 }
957 if (mDetectorInstances[detID]) {
958 LOG(error) << "O2HitMerger: DetID " << DetID::getName(detID) << " requested by external detector " << name
959 << " is already occupied; its hits will not be persisted. Assign a free DetID.";
960 delete extdet;
961 continue;
962 }
963 mDetectorInstances[detID].reset(extdet);
964 mExternalDetIDs.emplace_back(detID);
965 LOG(info) << "O2HitMerger: registered external detector " << name << " on DetID " << DetID::getName(detID)
966 << " (branch " << name << "Hit)";
967 }
968}
969
970} // namespace devices
971} // namespace o2
Definition of the DescriptorInnerBarrelITS3 class.
Definition of the Names Generator class.
Definition of the Stack class.
Sensitive detector built from an externally provided (CAD-derived) geometry.
Definition of the Detector class.
Definition of the Detector class.
Definition of the FV0 detector class.
std::vector< std::string > header
std::vector< long > entries
int32_t i
Definition of the Detector class.
Definition of the Detector class.
std::vector< o2::its::TrackITS > tracks
bool waitForControlInput(int workerID)
TBranch * ptr
ECal geometry creation and hit processing.
Definition of the Detector class.
Definition of the Detector class.
Definition of the Detector class.
StringRef key
static std::string getHitsFileName(DId d, const std::string_view prefix=STANDARDSIMPREFIX)
static int offsetTrackIndex(int trackID, int nprimaries, int primaryOffset, int secondaryOffset)
Definition Detector.h:178
static std::string getMCKinematicsFileName(const std::string_view prefix=STANDARDSIMPREFIX)
Definition NameConf.h:46
static std::string getMCHeadersFileName(const std::string_view prefix=STANDARDSIMPREFIX)
Definition NameConf.h:52
bool writeToDisc() const
Definition SimConfig.h:180
bool asService() const
Definition SimConfig.h:174
unsigned int getNEvents() const
Definition SimConfig.h:157
bool forwardKine() const
Definition SimConfig.h:179
static SimConfig & Instance()
Definition SimConfig.h:112
std::vector< std::string > const & getReadoutDetectors() const
Definition SimConfig.h:141
void putInfo(std::string const &key, T const &value)
void add(MCEventStats const &other)
merge from another object
Static class with identifiers, bitmasks and names for ALICE detectors.
Definition DetID.h:60
static constexpr const char * getName(ID id)
names of defined detectors
Definition DetID.h:148
static constexpr ID Last
if extra detectors added, update this !!!
Definition DetID.h:95
O2HitMerger()
Default constructor.
~O2HitMerger() override
Default destructor.
static ShmManager & Instance()
Definition ShmManager.h:61
const GLfloat * m
Definition glcorearb.h:4066
GLuint buffer
Definition glcorearb.h:655
GLuint entry
Definition glcorearb.h:5735
GLsizeiptr size
Definition glcorearb.h:659
GLuint index
Definition glcorearb.h:781
GLuint const GLchar * name
Definition glcorearb.h:781
GLsizei const GLfloat * value
Definition glcorearb.h:819
GLenum target
Definition glcorearb.h:1641
GLboolean * data
Definition glcorearb.h:298
GLint ref
Definition glcorearb.h:291
GLuint id
Definition glcorearb.h:650
uint8_t int statusCode int
TBranch * getOrMakeBranch(TTree &tree, const char *brname, T *ptr)
Definition Detector.h:310
bool parseSimReconfigFromString(std::string const &argumentstring, SimReconfigData &config)
int32_t const char * file
auto get(const std::byte *buffer, size_t=0)
Definition DataHeader.h:454
const bool const int TrackITSInternal< NLayers > & track
std::string simStatusString(std::string const &origin, std::string const &topic, std::string const &message)
bool publishMessage(fair::mq::Channel &channel, std::string const &message)
a couple of static helper functions to create timestamp values for CCDB queries or override obsolete ...
bool querySimConfig(fair::mq::Channel &channel)
TODO: Make this a base class of SimConfigData?
Definition SimConfig.h:205
LOG(info)<< "Compressed in "<< sw.CpuTime()<< " s"
std::unique_ptr< TTree > tree((TTree *) flIn.Get(std::string(o2::base::NameConf::CTFTREENAME).c_str()))