Project
Loading...
Searching...
No Matches
O2PrimaryServerDevice.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
15#include <fairmq/Device.h>
16#include <fairmq/TransportFactory.h>
17#include <FairPrimaryGenerator.h>
19#include <fairmq/Message.h>
20#include <DetectorsBase/Stack.h>
23#include <TMessage.h>
24#include <TClass.h>
29#include <SimConfig/SimConfig.h>
33#include <Field/MagneticField.h>
34#include <TGeoGlobalMagField.h>
35#include <typeinfo>
36#include <thread>
37#include <TROOT.h>
38#include <TStopwatch.h>
39#include <fstream>
40#include <iostream>
41#include <atomic>
42#include "PrimaryServerState.h"
44#include <chrono>
46#include <TRandom3.h>
47#include <regex>
48
49namespace o2
50{
51namespace devices
52{
53
55{
56 mUseFixedChunkSeed = getenv("ALICEO2_O2SIM_SUBEVENTSEED") && atoi(getenv("ALICEO2_O2SIM_SUBEVENTSEED"));
57 if (mUseFixedChunkSeed) {
58 mFixedChunkSeed = atol(getenv("ALICEO2_O2SIM_SUBEVENTSEED"));
59 }
60}
61
63{
64 try {
65 if (mGeneratorThread.joinable()) {
66 mGeneratorThread.join();
67 }
68 if (mControlThread.joinable()) {
69 mControlThread.join();
70 }
71 } catch (...) {
72 }
73}
74
76{
77 TStopwatch timer;
78 timer.Start();
79 const auto& conf = mSimConfig;
81 ccdbmgr.setURL(conf.getConfigData().mCCDBUrl);
82 ccdbmgr.setTimestamp(conf.getTimestamp());
83
84 // set the global information about the number of events to be generated
85 unsigned int nTotalEvents = conf.getNEvents();
87
88 // init magnetic field as it might be needed by the generator
89 if (TGeoGlobalMagField::Instance()->GetField() == nullptr) {
90 TGeoGlobalMagField::Instance()->SetField(o2::base::SimFieldUtils::createMagField());
91 TGeoGlobalMagField::Instance()->Lock();
92 }
93
94 // look if we find a cached instances of Pythia8 or external generators in order to avoid
95 // (long) initialization times.
96 // This is evidently a bit weak, as generators might need reconfiguration (to be treated later).
97 // For now, we'd like to allow for fast switches between say a pythia8 instance and reading from kinematics
98 // to continue an already started simulation.
99 //
100 // Not using cached instances for external kinematics since these might change input filenames etc.
101 // and are in any case quickly setup.
102 mPrimGen = nullptr;
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();
108 }
109 }
110
111 if (mPrimGen == nullptr) {
112 mPrimGen = new o2::eventgen::PrimaryGenerator;
114
115 // setup vertexing
116 auto vtxMode = conf.getVertexMode();
118 if (vtxMode == VertexMode::kNoVertex || vtxMode == VertexMode::kDiamondParam) {
119 mPrimGen->setVertexMode(vtxMode);
120 } else if (vtxMode == VertexMode::kCCDB) {
121 // we need to fetch the CCDB object
122 mPrimGen->setVertexMode(vtxMode, ccdbmgr.getForTimeStamp<o2::dataformats::MeanVertexObject>("GLO/Calib/MeanVertex", conf.getTimestamp()));
123 } else if (vtxMode == VertexMode::kCollCxt) {
124 // The vertex will be injected from the outside via setExternalVertex
125 } else {
126 LOG(fatal) << "Unsupported vertex mode";
127 }
128
129 auto embedinto_filename = conf.getEmbedIntoFileName();
130 if (!embedinto_filename.empty()) {
131 // determine the sim prefix from the embedding filename
132 // the filename should be an MCHeader file ... so it should match SOME_PATH/prefix_MCHeader.root
133 std::regex re(R"((.*/)?([^/]+)_MCHeader\.root$)");
134 std::smatch match;
135
136 if (std::regex_search(embedinto_filename, match, re)) {
137 std::cout << "Extracted embedding prefix : " << match[2] << '\n';
138 mEmbeddIntoPrefix = match[2];
139 } else {
140 LOG(fatal) << "Embedding asked but no suitable embedding prefix extractable from " << embedinto_filename;
141 }
142 mPrimGen->embedInto(embedinto_filename);
143 }
144
145 mPrimGen->Init();
146
147 std::unique_ptr<o2::eventgen::PrimaryGenerator> ptr_wrapper;
148 ptr_wrapper.reset(mPrimGen);
149 mPrimGeneratorCache[conf.getGenerator()] = std::move(ptr_wrapper);
150 }
151 mPrimGen->SetEvent(&mEventHeader);
152
153 // A good moment to couple to collision context
154 auto collContextFileName_PrefixPair = mSimConfig.getCollContextFilenameAndEventPrefix();
155 auto collContextFileName = collContextFileName_PrefixPair.first;
156 if (collContextFileName.size() > 0) {
157 LOG(info) << "Simulation has collission context";
158 mCollissionContext = o2::steer::DigitizationContext::loadFromFile(collContextFileName);
159 if (mCollissionContext) {
160 const auto& vertices = mCollissionContext->getInteractionVertices();
161 LOG(info) << "We found " << vertices.size() << " vertices included ";
162
163 // initialize the eventID to collID mapping
164 const auto source = mCollissionContext->findSimPrefix(collContextFileName_PrefixPair.second);
165 if (source == -1) {
166 LOG(fatal) << "Wrong simulation prefix";
167 }
168 mEventID_to_CollID.clear();
169 mEventID_to_CollID = mCollissionContext->getCollisionIndicesForSource(source);
170 }
171 }
172
173 LOG(info) << "Generator initialization took " << timer.CpuTime() << "s";
174 if (mMaxEvents > 0) {
175 generateEvent(); // generate a first event
176 }
177}
178
180{
181 bool changeState = true; // false;
182 LOG(info) << "Event generation started ";
183 if (changeState) {
185 }
186 TStopwatch timer;
187 timer.Start();
188 try {
189 bool valid = false;
190 int retry_counter = 0;
191 const int MAX_RETRY = 100;
192 do {
193 mStack->Reset();
194 const auto& conf = mSimConfig;
195 // see if we the vertex comes from the collision context
196 if (mCollissionContext && conf.getVertexMode() == o2::conf::VertexMode::kCollCxt) {
197 const auto& vertices = mCollissionContext->getInteractionVertices();
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";
202 mPrimGen->setExternalVertexForNextEvent(vertex.X(), vertex.Y(), vertex.Z());
203
204 // set correct embedding index for PrimaryGenerator ... based on collision context for embedding
205 auto& collisionParts = mCollissionContext->getEventParts()[collisionindex];
206 int background_index = -1; // -1 means no embedding taking place for this signal
207
208 // find the part that corresponds to the event embeded into
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;
213 }
214 }
215 mPrimGen->setEmbedIndex(background_index);
216 }
217 }
218 mPrimGen->GenerateEvent(mStack);
219 if (mStack->getPrimaries().size() > 0) {
220 valid = true;
221 } else {
222 retry_counter++;
223 if (retry_counter > MAX_RETRY) {
224 LOG(warn) << "Not able to generate a non-empty event in " << MAX_RETRY << " trials";
225 // empty event is sent out
226 valid = true;
227 }
228 }
229 } while (!valid);
230 } catch (std::exception const& e) {
231 LOG(error) << " Exception occurred during event gen " << e.what();
232 }
233 timer.Stop();
234 LOG(info) << "Event generation took " << timer.CpuTime() << "s"
235 << " and produced " << mStack->getPrimaries().size() << " primaries ";
236 if (changeState) {
238 }
239}
240
242{
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)));
247 channel.Send(reply);
248 };
249
250 LOG(info) << "LAUNCHING STATUS THREAD";
251 auto lambda = [this, sendErrorReply]() {
252 bool canShutdown{false};
253 // Exit only when both: serving stopped and allowed from outside.
254 while (!(mState == O2PrimaryServerState::Stopped && canShutdown)) {
255 auto& channel = GetChannels().at("o2sim-primserv-info").at(0);
256 if (!channel.IsValid()) {
257 LOG(error) << "channel primserv-info not valid";
258 }
259 std::unique_ptr<fair::mq::Message> request(channel.NewSimpleMessage((int)(-1)));
260 int timeout = 100; // 100ms --> so as not to block and allow for proper termination of this thread
261 if (channel.Receive(request, timeout) > 0) {
262 int request_payload; // we expect an (int) ~ to type O2PrimaryServerInfoRequest
263 if (request->GetSize() != sizeof(request_payload)) {
264 LOG(error) << "Obtained request with unexpected payload size";
265 sendErrorReply(channel); // ALWAYS reply
266 continue;
267 }
268
269 memcpy(&request_payload, request->GetData(), sizeof(request_payload));
270
271 if (request_payload == (int)O2PrimaryServerInfoRequest::Status) {
272 LOG(info) << "Received status request";
273 // request needs to be a simple enum of type O2PrimaryServerInfoRequest
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";
277 }
278 } else if (request_payload == (int)O2PrimaryServerInfoRequest::Config) {
279 HandleConfigRequest(channel);
280 } else if (request_payload == (int)O2PrimaryServerInfoRequest::AllowShutdown) {
281 LOG(info) << "Got info that we may shutdown";
282 std::unique_ptr<fair::mq::Message> ack(channel.NewSimpleMessage(200));
283 channel.Send(ack);
284 canShutdown = true;
285 } else {
286 sendErrorReply(channel);
287 }
288 }
289 }
290 mInfoThreadStopped = true;
291 };
292 threads.push_back(std::thread(lambda));
293 threads.back().detach();
294}
295
297{
298 // fatal without core dump
299 fair::Logger::OnFatal([] { throw fair::FatalException("Fatal error occured. Exiting without core dump..."); });
300
301 o2::simpubsub::publishMessage(GetChannels()["primary-notifications"].at(0), "SERVER : INITIALIZING");
302
304 LOG(info) << "Init Server device ";
305
306 // init sim config
307 auto& vm = GetConfig()->GetVarMap();
308 auto& conf = o2::conf::SimConfig::Instance();
309 if (vm.count("isRun5")) {
310 conf.setRun5();
311 }
312 conf.resetFromParsedMap(vm);
313
314 // update the parameters from an INI/JSON file, if given (overrides code-based version)
316 // update the parameters from stuff given at command line (overrides file-based version)
317 o2::conf::ConfigurableParam::updateFromString(conf.getKeyValueString());
318
319 // customize the level of log output
320 FairLogger::GetLogger()->SetLogScreenLevel(conf.getLogSeverity().c_str());
321 FairLogger::GetLogger()->SetLogVerbosityLevel(conf.getLogVerbosity().c_str());
322
323 // from now on mSimConfig should be used within this process
324 mSimConfig = conf;
325
326 mStack = new o2::data::Stack();
327 mStack->setExternalMode(true);
328
329 // MC ENGINE
330 LOG(info) << "ENGINE SET TO " << vm["mcEngine"].as<std::string>();
331 // CHUNK SIZE
332 mChunkGranularity = vm["chunkSize"].as<unsigned int>();
333 LOG(info) << "CHUNK SIZE SET TO " << mChunkGranularity;
334
335 // initial initial seed --> we should store this somewhere
336 mInitialSeed = vm["seed"].as<ULong_t>();
337 mInitialSeed = o2::utils::RngHelper::setGRandomSeed(mInitialSeed);
338 mSeedGenerator.SetSeed(mInitialSeed);
339 LOG(info) << "RNG INITIAL SEED " << mInitialSeed;
340
341 mMaxEvents = conf.getNEvents();
342
343 // need to make ROOT thread-safe since we use ROOT services in all places
344 ROOT::EnableThreadSafety();
345
347
348 // launch initialization of particle generator asynchronously
349 // so that we reach the RUNNING state of the server quickly
350 // and do not block here
351 mGeneratorThread = std::thread(&O2PrimaryServerDevice::initGenerator, this);
352 if (mGeneratorThread.joinable()) {
353 try {
354 mGeneratorThread.join();
355 } catch (std::exception const& e) {
356 LOG(warn) << "Exception during thread join ..ignoring";
357 }
358 }
359
360 // init pipe
361 auto pipeenv = getenv("ALICE_O2SIMSERVERTODRIVER_PIPE");
362 if (pipeenv) {
363 mPipeToDriver = atoi(pipeenv);
364 LOG(info) << "ASSIGNED PIPE HANDLE " << mPipeToDriver;
365 } else {
366 LOG(info) << "DID NOT FIND ENVIRONMENT VARIABLE TO INIT PIPE";
367 }
368
369 mAsService = vm["asservice"].as<bool>();
370 if (mAsService) {
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";
375 }
376 mControlChannel.Connect(std::string(controlsocketname));
377 mControlChannel.Validate();
378 }
379
380 if (mMaxEvents <= 0) {
381 if (mAsService) {
383 }
384 } else {
386 }
387
388 // feedback to driver that we are done initializing
389 if (mPipeToDriver != -1) {
390 int message = -111; // special code meaning end of initialization
391 if (write(mPipeToDriver, &message, sizeof(int))) {
392 }
393 }
394}
395
397{
398 LOG(info) << "ReInit Server device ";
399
400 if (reconfig.stop) {
401 return false;
402 }
403
404 // mSimConfig.getConfigData().mKeyValueTokens=reconfig.keyValueTokens;
405 // Think about this:
406 // update the parameters from an INI/JSON file, if given (overrides code-based version)
408 // update the parameters from stuff given at command line (overrides file-based version)
410
411 // initial initial seed --> we should store this somewhere
412 mInitialSeed = reconfig.startSeed;
413 mInitialSeed = o2::utils::RngHelper::setGRandomSeed(mInitialSeed);
414 mSeedGenerator.SetSeed(mInitialSeed);
415 LOG(info) << "RNG INITIAL SEED " << mInitialSeed;
416
417 mMaxEvents = reconfig.nEvents;
418
419 // updating the simconfig member with new information especially concerning the generators
420 // TODO: put this into utility function?
421 mSimConfig.getConfigData().mGenerator = reconfig.generator;
422 mSimConfig.getConfigData().mTrigger = reconfig.trigger;
423 mSimConfig.getConfigData().mExtKinFileName = reconfig.extKinfileName;
424
425 mEventCounter = 0;
426 mPartCounter = 0;
427 mNeedNewEvent = true;
428 // reinit generator and start generation of a new event
429 if (mGeneratorThread.joinable()) {
430 try {
431 mGeneratorThread.join();
432 } catch (std::exception const& e) {
433 LOG(warn) << "Exception during thread join ..ignoring";
434 }
435 }
436 // mGeneratorThread = std::thread(&O2PrimaryServerDevice::initGenerator, this);
438
439 return true;
440}
441
442bool O2PrimaryServerDevice::HandleConfigRequest(fair::mq::Channel& channel)
443{
444 LOG(info) << "Received config request";
445 // just sending the simulation configuration to anyone that wants it
446 const auto& confdata = mSimConfig.getConfigData();
447
448 TMessage* tmsg = new TMessage(kMESS_OBJECT);
449 tmsg->WriteObjectAny((void*)&confdata, TClass::GetClass(typeid(confdata)));
450
451 auto free_tmessage = [](void* data, void* hint) { delete static_cast<TMessage*>(hint); };
452
453 std::unique_ptr<fair::mq::Message> message(
454 fTransportFactory->CreateMessage(tmsg->Buffer(), tmsg->BufferSize(), free_tmessage, tmsg));
455
456 // send answer
457 if (channel.Send(message) > 0) {
458 LOG(info) << "config reply send ";
459 return true;
460 } else {
461 LOG(error) << "Failure sending config reply ";
462 }
463 return true;
464}
465
467{
468 // we might come here in IDLE mode
469 if (mState.load() == O2PrimaryServerState::Idle) {
470 if (mWaitingControlInput.load() == 0) {
471 if (mControlThread.joinable()) {
472 mControlThread.join();
473 }
474 mControlThread = std::thread(&O2PrimaryServerDevice::waitForControlInput, this);
475 }
476 }
477
478 auto& channel = GetChannels().at("primary-get").at(0);
479 PrimaryChunkRequest requestpayload;
480 std::unique_ptr<fair::mq::Message> request(channel.NewSimpleMessage(requestpayload));
481 auto bytes = channel.Receive(request);
482 if (bytes < 0) {
483 LOG(error) << "Some error/interrupt occurred on socket during receive";
484 if (NewStatePending()) { // new state is typically pending if (term) signal was received
485 WaitForNextState();
486 // ask ourselves for termination of this loop
488 }
489 return false;
490 }
491
492 TStopwatch timer;
493 timer.Start();
494 auto& r = *((PrimaryChunkRequest*)(request->GetData()));
495 LOG(debug) << "PARTICLE REQUEST IN STATE " << PrimStateToString[(int)mState.load()] << " from " << r.workerid << ":" << r.requestid;
496
497 auto prestate = mState.load();
498 auto more = HandleRequest(request, 0, channel);
499 if (!more) {
500 if (mAsService) {
503 }
504 } else {
506 }
507 }
508 timer.Stop();
509 auto time = timer.CpuTime();
510 LOG(debug) << "COND-RUN TOOK " << time << " s";
511 return mState != O2PrimaryServerState::Stopped;
512}
513
515{
516 // We shouldn't shut down immediately when all events have been served
517 // Instead we also need to wait until the info thread running some communication server
518 // with other processes is finished.
519 while (!mInfoThreadStopped) {
520 LOG(info) << "Waiting info thread";
521 using namespace std::chrono_literals;
522 std::this_thread::sleep_for(1000ms);
523 }
524}
525
526bool O2PrimaryServerDevice::HandleRequest(fair::mq::MessagePtr& request, int /*index*/, fair::mq::Channel& channel)
527{
528 // LOG(debug) << "GOT A REQUEST WITH SIZE " << request->GetSize();
529 // std::string requeststring(static_cast<char*>(request->GetData()), request->GetSize());
530 // LOG(info) << "NORMAL REQUEST STRING " << requeststring;
531 bool workavailable = true;
532 if (mEventCounter >= mMaxEvents && mNeedNewEvent) {
533 workavailable = false;
534 }
535 if (!(mState.load() == O2PrimaryServerState::ReadyToServe || mState.load() == O2PrimaryServerState::WaitingEvent)) {
536 // send a zero answer
537 workavailable = false;
538 }
539
540 PrimaryChunkAnswer header{mState, workavailable};
541 fair::mq::Parts reply;
542 std::unique_ptr<fair::mq::Message> headermsg(channel.NewSimpleMessage(header));
543 reply.AddPart(std::move(headermsg));
544
545 LOG(debug) << "Received request for work " << mEventCounter << " " << mMaxEvents << " " << mNeedNewEvent << " available " << workavailable;
546 if (workavailable) {
547
548 if (mNeedNewEvent) {
549 // we need a newly generated event now
550 if (mGeneratorThread.joinable()) {
551 try {
552 mGeneratorThread.join();
553 } catch (std::exception const& e) {
554 LOG(warn) << "Exception during thread join ..ignoring";
555 }
556 }
557 // also if we are still in event waiting stage (doing some busy sleep)
558 while (mState.load() == O2PrimaryServerState::WaitingEvent) {
559 LOG(info) << "Waiting for event generation do become fully available";
560 usleep(100);
561 }
562 mNeedNewEvent = false;
563 mPartCounter = 0;
564 mEventCounter++;
565 }
566
567 auto& prims = mStack->getPrimaries();
568 auto numberofparts = (int)std::ceil(prims.size() / (1. * mChunkGranularity));
569 // number of parts should be at least 1 (even if empty)
570 numberofparts = std::max(1, numberofparts);
571
572 LOG(debug) << "Have " << prims.size() << " " << numberofparts;
573
576 i.eventID = workavailable ? mEventCounter : -1;
577 i.maxEvents = mMaxEvents;
578 i.part = mPartCounter + 1;
579 i.nparts = numberofparts;
580 // assign a deterministic (yet collision free seed) to process this particle chunk in Geant
581 // limit range to uint32_t since internal limit of TRandom (despite API suggesting otherwise)
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;
586 m.mSubEventInfo = i;
587
588 int endindex = prims.size() - mPartCounter * mChunkGranularity;
589 int startindex = prims.size() - (mPartCounter + 1) * mChunkGranularity;
590 LOG(debug) << "indices " << startindex << " " << endindex;
591
592 if (startindex < 0) {
593 startindex = 0;
594 }
595 if (endindex < 0) {
596 endindex = 0;
597 }
598
599 for (int index = startindex; index < endindex; ++index) {
600 m.mParticles.emplace_back(prims[index]);
601 }
602
603 LOG(info) << "Sending " << m.mParticles.size() << " particles";
604 LOG(info) << "treating ev " << mEventCounter << " part " << i.part << " out of " << i.nparts;
605
606 // feedback to driver if new event started
607 if (mPipeToDriver != -1 && i.part == 1 && workavailable) {
608 if (write(mPipeToDriver, &mEventCounter, sizeof(mEventCounter))) {
609 }
610 }
611
612 mPartCounter++;
613 if (mPartCounter == numberofparts) {
614 mNeedNewEvent = true;
615 // start generation of a new event
616 if (mEventCounter < mMaxEvents) {
617 mGeneratorThread = std::thread(&O2PrimaryServerDevice::generateEvent, this);
618 }
619 }
620
621 TMessage* tmsg = new TMessage(kMESS_OBJECT);
622 tmsg->WriteObjectAny((void*)&m, TClass::GetClass("o2::data::PrimaryChunk"));
623
624 auto free_tmessage = [](void* data, void* hint) { delete static_cast<TMessage*>(hint); };
625
626 std::unique_ptr<fair::mq::Message> message(channel.NewMessage(tmsg->Buffer(), tmsg->BufferSize(), free_tmessage, tmsg));
627
628 reply.AddPart(std::move(message));
629 }
630
631 // send answer
632 TStopwatch timer;
633 timer.Start();
634 auto code = Send(reply, "primary-get", 0, 5000); // we introduce timeout in order not to block other requests
635 timer.Stop();
636 auto time = timer.CpuTime();
637 if (code > 0) {
638 LOG(debug) << "Reply send in " << time << "s";
639 return workavailable;
640 } else {
641 LOG(warn) << "Sending process had problems. Return code : " << code << " time " << time << "s";
642 }
643 return false; // -> error should not get here
644}
645
647{
648 LOG(info) << message << " CHANGING STATE TO " << PrimStateToString[(int)to];
649 mState = to;
650}
651
653{
654 mWaitingControlInput.store(1);
655 if (mState.load() != O2PrimaryServerState::Idle) {
656 mWaitingControlInput.store(0);
657 return;
658 }
659
660 o2::simpubsub::publishMessage(GetChannels()["primary-notifications"].at(0), o2::simpubsub::simStatusString("PRIMSERVER", "STATUS", "AWAITING INPUT"));
661 // this means we are idling
662
663 std::unique_ptr<fair::mq::Message> reply(mControlChannel.NewMessage());
664
665 bool ok = false;
666
667 LOG(info) << "WAITING FOR CONTROL INPUT";
668 if (mControlChannel.Receive(reply) > 0) {
670 auto data = reply->GetData();
671 auto size = reply->GetSize();
672
673 std::string command(reinterpret_cast<char const*>(data), size);
674 LOG(info) << "message: " << command;
675
677 o2::conf::parseSimReconfigFromString(command, reconfig);
678 LOG(info) << "Processing " << reconfig.nEvents << " new events";
679 try {
680 LOG(info) << "REINIT START";
681 ok = ReInit(reconfig);
682 LOG(info) << "REINIT DONE";
683 } catch (std::exception e) {
684 LOG(info) << "Exception during reinit";
685 }
686 } else {
687 LOG(info) << "NOTHING RECEIVED";
688 }
689 if (ok) {
690 // stateTransition(O2PrimaryServerState::ReadyToServe, "CONTROL"); --> SHOULD BE DONE FROM EVENT GENERATOR (which get's however called only when mEvents>0)
691 } else {
693 }
694 mWaitingControlInput.store(0);
695}
696
697} // namespace devices
698} // namespace o2
Definition of the Stack class.
std::ostringstream debug
uint64_t vertex
Definition RawEventData.h:9
int16_t time
Definition RawEventData.h:4
std::vector< std::string > header
int32_t i
Definition of the MagF class.
bool valid
Methods to create simulation mag field.
static FairField *const createMagField()
static BasicCCDBManager & instance()
static void updateFromFile(std::string const &, std::string const &paramsList="", bool unchangedOnly=false)
static void updateFromString(std::string const &)
SimConfigData const & getConfigData() const
Definition SimConfig.h:134
std::pair< std::string, std::string > getCollContextFilenameAndEventPrefix() const
static SimConfig & Instance()
Definition SimConfig.h:112
std::string getOutPrefix() const
Definition SimConfig.h:164
void setExternalMode(bool m)
Definition Stack.h:205
void Reset() override
Resets arrays and stack and deletes particles and tracks.
Definition Stack.cxx:588
const std::vector< TParticle > & getPrimaries() const
Definition Stack.h:187
bool ReInit(o2::conf::SimReconfigData const &reconfig)
void stateTransition(O2PrimaryServerState to, const char *message)
~O2PrimaryServerDevice() final
Default destructor.
bool HandleConfigRequest(fair::mq::Channel &channel)
bool HandleRequest(fair::mq::MessagePtr &request, int, fair::mq::Channel &channel)
static void setTotalNEvents(unsigned int &n)
Definition Generator.h:96
Bool_t GenerateEvent(FairGenericStack *pStack) override
void setVertexMode(o2::conf::VertexMode const &mode, o2::dataformats::MeanVertexObject const *obj=nullptr)
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)
Definition RngHelper.h:37
bool match(const std::vector< std::string > &queries, const char *pattern)
Definition dcs-ccdb.cxx:229
const GLfloat * m
Definition glcorearb.h:4066
GLsizeiptr size
Definition glcorearb.h:659
GLuint index
Definition glcorearb.h:781
GLsizei GLsizei GLchar * source
Definition glcorearb.h:798
GLboolean * data
Definition glcorearb.h:298
GLuint GLsizei const GLchar * message
Definition glcorearb.h:2517
GLboolean r
Definition glcorearb.h:1233
GLbitfield GLuint64 timeout
Definition glcorearb.h:1573
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
Definition SimConfig.h:58
std::string mGenerator
Definition SimConfig.h:55
TODO: Make this a base class of SimConfigData?
Definition SimConfig.h:205
static void setPrimaryGenerator(o2::conf::SimConfig const &, FairPrimaryGenerator *)
LOG(info)<< "Compressed in "<< sw.CpuTime()<< " s"