15#include <fairmq/Device.h>
16#include <fairmq/TransportFactory.h>
17#include <FairPrimaryGenerator.h>
19#include <fairmq/Message.h>
34#include <TGeoGlobalMagField.h>
38#include <TStopwatch.h>
56 mUseFixedChunkSeed = getenv(
"ALICEO2_O2SIM_SUBEVENTSEED") && atoi(getenv(
"ALICEO2_O2SIM_SUBEVENTSEED"));
57 if (mUseFixedChunkSeed) {
58 mFixedChunkSeed = atol(getenv(
"ALICEO2_O2SIM_SUBEVENTSEED"));
65 if (mGeneratorThread.joinable()) {
66 mGeneratorThread.join();
68 if (mControlThread.joinable()) {
69 mControlThread.join();
79 const auto& conf = mSimConfig;
81 ccdbmgr.setURL(conf.getConfigData().mCCDBUrl);
82 ccdbmgr.setTimestamp(conf.getTimestamp());
85 unsigned int nTotalEvents = conf.getNEvents();
89 if (TGeoGlobalMagField::Instance()->GetField() ==
nullptr) {
91 TGeoGlobalMagField::Instance()->Lock();
103 if (conf.getGenerator().compare(
"extkin") != 0 && conf.getGenerator().compare(
"extkinO2") != 0) {
104 auto iter = mPrimGeneratorCache.find(conf.getGenerator());
105 if (iter != mPrimGeneratorCache.end()) {
106 mPrimGen = iter->second.get();
107 LOG(info) <<
"Found cached generator for " << conf.getGenerator();
111 if (mPrimGen ==
nullptr) {
116 auto vtxMode = conf.getVertexMode();
118 if (vtxMode == VertexMode::kNoVertex || vtxMode == VertexMode::kDiamondParam) {
120 }
else if (vtxMode == VertexMode::kCCDB) {
123 }
else if (vtxMode == VertexMode::kCollCxt) {
126 LOG(fatal) <<
"Unsupported vertex mode";
129 auto embedinto_filename = conf.getEmbedIntoFileName();
130 if (!embedinto_filename.empty()) {
133 std::regex re(R
"((.*/)?([^/]+)_MCHeader\.root$)");
136 if (std::regex_search(embedinto_filename,
match, re)) {
137 std::cout <<
"Extracted embedding prefix : " <<
match[2] <<
'\n';
138 mEmbeddIntoPrefix =
match[2];
140 LOG(fatal) <<
"Embedding asked but no suitable embedding prefix extractable from " << embedinto_filename;
147 std::unique_ptr<o2::eventgen::PrimaryGenerator> ptr_wrapper;
148 ptr_wrapper.reset(mPrimGen);
149 mPrimGeneratorCache[conf.getGenerator()] = std::move(ptr_wrapper);
151 mPrimGen->SetEvent(&mEventHeader);
155 auto collContextFileName = collContextFileName_PrefixPair.first;
156 if (collContextFileName.size() > 0) {
157 LOG(info) <<
"Simulation has collission context";
159 if (mCollissionContext) {
161 LOG(info) <<
"We found " << vertices.size() <<
" vertices included ";
164 const auto source = mCollissionContext->
findSimPrefix(collContextFileName_PrefixPair.second);
166 LOG(fatal) <<
"Wrong simulation prefix";
168 mEventID_to_CollID.clear();
173 LOG(info) <<
"Generator initialization took " << timer.CpuTime() <<
"s";
174 if (mMaxEvents > 0) {
181 bool changeState =
true;
182 LOG(info) <<
"Event generation started ";
190 int retry_counter = 0;
191 const int MAX_RETRY = 100;
194 const auto& conf = mSimConfig;
198 if (vertices.size() > 0) {
199 auto collisionindex = mEventID_to_CollID.at(mEventCounter);
200 auto&
vertex = vertices.at(collisionindex);
201 LOG(info) <<
"Setting vertex " <<
vertex <<
" for event " << mEventCounter <<
" for prefix " << mSimConfig.
getOutPrefix() <<
" from CollContext";
205 auto& collisionParts = mCollissionContext->
getEventParts()[collisionindex];
206 int background_index = -1;
209 for (
auto& part : collisionParts) {
210 if (mCollissionContext->
getSimPrefixes()[part.sourceID] == mEmbeddIntoPrefix) {
211 background_index = part.entryID;
212 LOG(info) <<
"Setting embedding index to " << background_index;
223 if (retry_counter > MAX_RETRY) {
224 LOG(warn) <<
"Not able to generate a non-empty event in " << MAX_RETRY <<
" trials";
230 }
catch (std::exception
const& e) {
231 LOG(error) <<
" Exception occurred during event gen " << e.what();
234 LOG(info) <<
"Event generation took " << timer.CpuTime() <<
"s"
235 <<
" and produced " << mStack->
getPrimaries().size() <<
" primaries ";
243 static std::vector<std::thread> threads;
244 auto sendErrorReply = [](fair::mq::Channel& channel) {
245 LOG(error) <<
"UNKNOWN REQUEST";
246 std::unique_ptr<fair::mq::Message> reply(channel.NewSimpleMessage((
int)(404)));
250 LOG(info) <<
"LAUNCHING STATUS THREAD";
251 auto lambda = [
this, sendErrorReply]() {
252 bool canShutdown{
false};
255 auto& channel = GetChannels().at(
"o2sim-primserv-info").at(0);
256 if (!channel.IsValid()) {
257 LOG(error) <<
"channel primserv-info not valid";
259 std::unique_ptr<fair::mq::Message> request(channel.NewSimpleMessage((
int)(-1)));
261 if (channel.Receive(request,
timeout) > 0) {
263 if (request->GetSize() !=
sizeof(request_payload)) {
264 LOG(error) <<
"Obtained request with unexpected payload size";
265 sendErrorReply(channel);
269 memcpy(&request_payload, request->GetData(),
sizeof(request_payload));
272 LOG(info) <<
"Received status request";
274 std::unique_ptr<fair::mq::Message> reply(channel.NewSimpleMessage((
int)mState.load()));
275 if (channel.Send(reply) > 0) {
276 LOG(info) <<
"Send status successful";
281 LOG(info) <<
"Got info that we may shutdown";
282 std::unique_ptr<fair::mq::Message> ack(channel.NewSimpleMessage(200));
286 sendErrorReply(channel);
290 mInfoThreadStopped =
true;
292 threads.push_back(std::thread(lambda));
293 threads.back().detach();
299 fair::Logger::OnFatal([] {
throw fair::FatalException(
"Fatal error occured. Exiting without core dump..."); });
304 LOG(info) <<
"Init Server device ";
307 auto& vm = GetConfig()->GetVarMap();
309 if (vm.count(
"isRun5")) {
312 conf.resetFromParsedMap(vm);
320 FairLogger::GetLogger()->SetLogScreenLevel(conf.getLogSeverity().c_str());
321 FairLogger::GetLogger()->SetLogVerbosityLevel(conf.getLogVerbosity().c_str());
330 LOG(info) <<
"ENGINE SET TO " << vm[
"mcEngine"].as<std::string>();
332 mChunkGranularity = vm[
"chunkSize"].as<
unsigned int>();
333 LOG(info) <<
"CHUNK SIZE SET TO " << mChunkGranularity;
336 mInitialSeed = vm[
"seed"].as<ULong_t>();
338 mSeedGenerator.SetSeed(mInitialSeed);
339 LOG(info) <<
"RNG INITIAL SEED " << mInitialSeed;
341 mMaxEvents = conf.getNEvents();
344 ROOT::EnableThreadSafety();
352 if (mGeneratorThread.joinable()) {
354 mGeneratorThread.join();
355 }
catch (std::exception
const& e) {
356 LOG(warn) <<
"Exception during thread join ..ignoring";
361 auto pipeenv = getenv(
"ALICE_O2SIMSERVERTODRIVER_PIPE");
363 mPipeToDriver = atoi(pipeenv);
364 LOG(info) <<
"ASSIGNED PIPE HANDLE " << mPipeToDriver;
366 LOG(info) <<
"DID NOT FIND ENVIRONMENT VARIABLE TO INIT PIPE";
369 mAsService = vm[
"asservice"].as<
bool>();
371 mControlChannel = fair::mq::Channel{
"o2sim-control",
"sub", fTransportFactory};
372 auto controlsocketname = getenv(
"ALICE_O2SIMCONTROL");
373 if (!controlsocketname) {
374 LOG(fatal) <<
"Internal error: Socketname for control input missing";
376 mControlChannel.Connect(std::string(controlsocketname));
377 mControlChannel.Validate();
380 if (mMaxEvents <= 0) {
389 if (mPipeToDriver != -1) {
391 if (write(mPipeToDriver, &
message,
sizeof(
int))) {
398 LOG(info) <<
"ReInit Server device ";
414 mSeedGenerator.SetSeed(mInitialSeed);
415 LOG(info) <<
"RNG INITIAL SEED " << mInitialSeed;
427 mNeedNewEvent =
true;
429 if (mGeneratorThread.joinable()) {
431 mGeneratorThread.join();
432 }
catch (std::exception
const& e) {
433 LOG(warn) <<
"Exception during thread join ..ignoring";
444 LOG(info) <<
"Received config request";
449 tmsg->WriteObjectAny((
void*)&confdata, TClass::GetClass(
typeid(confdata)));
451 auto free_tmessage = [](
void*
data,
void* hint) {
delete static_cast<TMessage*
>(hint); };
453 std::unique_ptr<fair::mq::Message>
message(
454 fTransportFactory->CreateMessage(tmsg->Buffer(), tmsg->BufferSize(), free_tmessage, tmsg));
457 if (channel.Send(
message) > 0) {
458 LOG(info) <<
"config reply send ";
461 LOG(error) <<
"Failure sending config reply ";
470 if (mWaitingControlInput.load() == 0) {
471 if (mControlThread.joinable()) {
472 mControlThread.join();
478 auto& channel = GetChannels().at(
"primary-get").at(0);
480 std::unique_ptr<fair::mq::Message> request(channel.NewSimpleMessage(requestpayload));
481 auto bytes = channel.Receive(request);
483 LOG(error) <<
"Some error/interrupt occurred on socket during receive";
484 if (NewStatePending()) {
495 LOG(
debug) <<
"PARTICLE REQUEST IN STATE " <<
PrimStateToString[(int)mState.load()] <<
" from " <<
r.workerid <<
":" <<
r.requestid;
497 auto prestate = mState.load();
509 auto time = timer.CpuTime();
519 while (!mInfoThreadStopped) {
520 LOG(info) <<
"Waiting info thread";
521 using namespace std::chrono_literals;
522 std::this_thread::sleep_for(1000ms);
531 bool workavailable =
true;
532 if (mEventCounter >= mMaxEvents && mNeedNewEvent) {
533 workavailable =
false;
537 workavailable =
false;
541 fair::mq::Parts reply;
542 std::unique_ptr<fair::mq::Message> headermsg(channel.NewSimpleMessage(
header));
543 reply.AddPart(std::move(headermsg));
545 LOG(
debug) <<
"Received request for work " << mEventCounter <<
" " << mMaxEvents <<
" " << mNeedNewEvent <<
" available " << workavailable;
550 if (mGeneratorThread.joinable()) {
552 mGeneratorThread.join();
553 }
catch (std::exception
const& e) {
554 LOG(warn) <<
"Exception during thread join ..ignoring";
559 LOG(info) <<
"Waiting for event generation do become fully available";
562 mNeedNewEvent =
false;
568 auto numberofparts = (int)std::ceil(prims.size() / (1. * mChunkGranularity));
570 numberofparts = std::max(1, numberofparts);
572 LOG(
debug) <<
"Have " << prims.size() <<
" " << numberofparts;
576 i.
eventID = workavailable ? mEventCounter : -1;
577 i.maxEvents = mMaxEvents;
578 i.part = mPartCounter + 1;
579 i.nparts = numberofparts;
582 const uint64_t drawnSeed = (uint64_t)(
static_cast<double>(std::numeric_limits<uint32_t>::max()) * mSeedGenerator.Rndm());
583 i.seed = mUseFixedChunkSeed ? mFixedChunkSeed : drawnSeed;
584 i.index =
m.mParticles.size();
585 i.mMCEventHeader = mEventHeader;
588 int endindex = prims.size() - mPartCounter * mChunkGranularity;
589 int startindex = prims.size() - (mPartCounter + 1) * mChunkGranularity;
590 LOG(
debug) <<
"indices " << startindex <<
" " << endindex;
592 if (startindex < 0) {
600 m.mParticles.emplace_back(prims[
index]);
603 LOG(info) <<
"Sending " <<
m.mParticles.size() <<
" particles";
604 LOG(info) <<
"treating ev " << mEventCounter <<
" part " <<
i.part <<
" out of " <<
i.nparts;
607 if (mPipeToDriver != -1 &&
i.part == 1 && workavailable) {
608 if (write(mPipeToDriver, &mEventCounter,
sizeof(mEventCounter))) {
613 if (mPartCounter == numberofparts) {
614 mNeedNewEvent =
true;
616 if (mEventCounter < mMaxEvents) {
622 tmsg->WriteObjectAny((
void*)&
m, TClass::GetClass(
"o2::data::PrimaryChunk"));
624 auto free_tmessage = [](
void*
data,
void* hint) {
delete static_cast<TMessage*
>(hint); };
626 std::unique_ptr<fair::mq::Message>
message(channel.NewMessage(tmsg->Buffer(), tmsg->BufferSize(), free_tmessage, tmsg));
628 reply.AddPart(std::move(
message));
634 auto code = Send(reply,
"primary-get", 0, 5000);
636 auto time = timer.CpuTime();
639 return workavailable;
641 LOG(warn) <<
"Sending process had problems. Return code : " << code <<
" time " <<
time <<
"s";
654 mWaitingControlInput.store(1);
656 mWaitingControlInput.store(0);
663 std::unique_ptr<fair::mq::Message> reply(mControlChannel.NewMessage());
667 LOG(info) <<
"WAITING FOR CONTROL INPUT";
668 if (mControlChannel.Receive(reply) > 0) {
670 auto data = reply->GetData();
671 auto size = reply->GetSize();
673 std::string command(
reinterpret_cast<char const*
>(
data),
size);
674 LOG(info) <<
"message: " << command;
678 LOG(info) <<
"Processing " << reconfig.
nEvents <<
" new events";
680 LOG(info) <<
"REINIT START";
682 LOG(info) <<
"REINIT DONE";
683 }
catch (std::exception e) {
684 LOG(info) <<
"Exception during reinit";
687 LOG(info) <<
"NOTHING RECEIVED";
694 mWaitingControlInput.store(0);
Definition of the Stack class.
std::vector< std::string > header
Definition of the MagF class.
Methods to create simulation mag field.
static FairField *const createMagField()
static BasicCCDBManager & instance()
static void updateFromFile(std::string const &, std::string const ¶msList="", bool unchangedOnly=false)
static void updateFromString(std::string const &)
SimConfigData const & getConfigData() const
std::pair< std::string, std::string > getCollContextFilenameAndEventPrefix() const
static SimConfig & Instance()
std::string getOutPrefix() const
void setExternalMode(bool m)
void Reset() override
Resets arrays and stack and deletes particles and tracks.
const std::vector< TParticle > & getPrimaries() const
bool ReInit(o2::conf::SimReconfigData const &reconfig)
bool ConditionalRun() override
void stateTransition(O2PrimaryServerState to, const char *message)
~O2PrimaryServerDevice() final
Default destructor.
bool HandleConfigRequest(fair::mq::Channel &channel)
void waitForControlInput()
bool HandleRequest(fair::mq::MessagePtr &request, int, fair::mq::Channel &channel)
O2PrimaryServerDevice()
constructor
static void setTotalNEvents(unsigned int &n)
Bool_t GenerateEvent(FairGenericStack *pStack) override
void setVertexMode(o2::conf::VertexMode const &mode, o2::dataformats::MeanVertexObject const *obj=nullptr)
Bool_t embedInto(TString fname)
void setEmbedIndex(int idx)
sets the embedding index
void setExternalVertexForNextEvent(double x, double y, double z)
std::unordered_map< int, int > getCollisionIndicesForSource(int source) const
std::vector< math_utils::Point3D< float > > const & getInteractionVertices() const
int findSimPrefix(std::string const &prefix) const
std::vector< std::vector< o2::steer::EventPart > > & getEventParts(bool withQED=false)
std::vector< std::string > const & getSimPrefixes() const
static DigitizationContext * loadFromFile(std::string_view filename="")
static ULong_t setGRandomSeed(ULong_t seed=0)
bool match(const std::vector< std::string > &queries, const char *pattern)
GLsizei GLsizei GLchar * source
GLuint GLsizei const GLchar * message
GLbitfield GLuint64 timeout
bool parseSimReconfigFromString(std::string const &argumentstring, SimReconfigData &config)
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 ...
constexpr const char * PrimStateToString[5]
O2PrimaryServerState
enum to represent state of the O2Sim event/primary server
std::string mExtKinFileName
TODO: Make this a base class of SimConfigData?
std::string extKinfileName
std::string keyValueTokens
static void setPrimaryGenerator(o2::conf::SimConfig const &, FairPrimaryGenerator *)
LOG(info)<< "Compressed in "<< sw.CpuTime()<< " s"