Project
Loading...
Searching...
No Matches
O2SimDevice.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 "O2SimDevice.h"
15#include "../macro/o2sim.C"
16#include "TVirtualMC.h"
17#include <fairmq/Message.h>
18#include <fairmq/Parts.h>
19#include <fairlogger/Logger.h>
20#include <DetectorsBase/Stack.h>
23#include <TRandom.h>
24#include <SimConfig/SimConfig.h>
25#include <cstring>
26
27void doLogInfo(int workerID, std::string const& message)
28{
29 LOG(info) << "[W" << workerID << "] " << message;
30}
31
32namespace o2
33{
34namespace devices
35{
36
38{
39 FairSystemInfo sysinfo;
41 LOG(info) << "Shutting down O2SimDevice";
42 LOG(info) << "TIME-STAMP " << mTimer.RealTime() << "\t";
43 LOG(info) << "MEM-STAMP " << sysinfo.GetCurrentMemory() / (1024. * 1024) << " " << sysinfo.GetMaxMemory() << " MB\n";
44}
45
47{
48 // in the initialization phase we will init the simulation
49 // NOTE: In a fair::mq::Device this is better done here (instead of outside) since
50 // we have to setup simulation + worker in the same thread (due to many threadlocal variables
51 // in the simulation) ... at least as long fair::mq::Device is not spawning workers on the master thread
52 initSim(GetChannels().at("o2sim-primserv-info").at(0), mSimRun);
53
54 // set the vmc and app pointers
55 mVMC = TVirtualMC::GetMC();
56 mVMCApp = static_cast<o2::steer::O2MCApplication*>(TVirtualMCApplication::Instance());
57 lateInit();
58}
59
61{
62 // late init
63 mVMCApp->initLate();
64}
65
66bool O2SimDevice::initSim(fair::mq::Channel& channel, std::unique_ptr<FairRunSim>& simptr)
67{
68 if (!o2::querySimConfig(channel)) {
69 return false;
70 }
71
72 LOG(info) << "Setting up the simulation ...";
73 simptr = std::move(std::unique_ptr<FairRunSim>(o2sim_init(true)));
74 FairSystemInfo sysinfo;
75
76 // to finish initialization (trigger further cross section table building etc) -- which especially
77 // G4 is doing at the first ProcessRun
78 // The goal is to have everything setup before we fork
79 TVirtualMC::GetMC()->ProcessRun(0);
80
81 LOG(info) << "MEM-STAMP END OF SIM INIT" << sysinfo.GetCurrentMemory() / (1024. * 1024) << " "
82 << sysinfo.GetMaxMemory() << " MB\n";
83
84 return true;
85}
86
87bool O2SimDevice::isWorkAvailable(fair::mq::Channel& statuschannel, int workerID)
88{
89 std::stringstream str;
90 str << "[W" << workerID << "]";
91 auto workerStr = str.str();
92
93 int timeoutinMS = 2000; // wait for 2s max
94 bool reprobe = true;
95 while (reprobe) {
96 reprobe = false;
97 int i = -1;
98 fair::mq::MessagePtr request(statuschannel.NewSimpleMessage((int)O2PrimaryServerInfoRequest::Status));
99 fair::mq::MessagePtr reply(statuschannel.NewSimpleMessage(i));
100 auto sendcode = statuschannel.Send(request, timeoutinMS);
101 if (sendcode > 0) {
102 LOG(info) << workerStr << " Waiting for status answer ";
103 auto code = statuschannel.Receive(reply, timeoutinMS);
104 if (code > 0) {
105 int state(*((int*)(reply->GetData())));
107 LOG(info) << workerStr << " SERVER IS SERVING";
108 return true;
110 LOG(info) << workerStr << " SERVER IS STILL INITIALIZING";
111 reprobe = true;
112 sleep(1);
114 LOG(info) << workerStr << " SERVER IS WAITING FOR EVENT";
115 reprobe = true;
116 sleep(1);
117 } else if (state == (int)o2::O2PrimaryServerState::Idle) {
118 LOG(info) << workerStr << " SERVER IS IDLE";
119 return false;
120 } else {
121 LOG(info) << workerStr << " SERVER STATE UNKNOWN OR STOPPED";
122 }
123 } else {
124 LOG(error) << workerStr << " STATUS REQUEST UNSUCCESSFUL";
125 }
126 }
127 }
128 return false;
129}
130
131bool O2SimDevice::Kernel(int workerID, fair::mq::Channel& requestchannel, fair::mq::Channel& dataoutchannel, fair::mq::Channel* statuschannel)
132{
133 static int counter = 0;
134 bool reproducibleSim = true;
135 if (getenv("O2_DISABLE_REPRODUCIBLE_SIM")) {
136 reproducibleSim = false;
137 }
138
139 // Mainly for debugging reasons, we allow to transport
140 // a specific event + eventpart. This allows to reproduce and debug bugs faster, once
141 // we know in which precise chunk they occur. The expected format for the environment variable
142 // is "eventnum:partid".
143 auto eventselection = getenv("O2SIM_RESTRICT_EVENTPART");
144 int focus_on_event = -1;
145 int focus_on_part = -1;
146 if (eventselection) {
147 auto splitString = [](const std::string& str) {
148 std::pair<std::string, std::string> parts;
149 size_t pos = str.find(':');
150 if (pos != std::string::npos) {
151 parts.first = str.substr(0, pos);
152 parts.second = str.substr(pos + 1);
153 }
154 return parts;
155 };
156 auto p = splitString(eventselection);
157 focus_on_event = std::atoi(p.first.c_str());
158 focus_on_part = std::atoi(p.second.c_str());
159 }
160
161 fair::mq::MessagePtr request(requestchannel.NewSimpleMessage(PrimaryChunkRequest{workerID, -1, counter++})); // <-- don't need content; channel means -> give primaries
162 fair::mq::Parts reply;
163
164 mVMCApp->setSimDataChannel(&dataoutchannel);
165
166 // we log info with workerID prepended
167 auto workerStr = [workerID]() {
168 std::stringstream str;
169 str << "[W" << workerID << "]";
170 return str.str();
171 };
172
173 doLogInfo(workerID, "Requesting work chunk");
174 int timeoutinMS = 2000;
175 auto sendcode = requestchannel.Send(request, timeoutinMS);
176 if (sendcode > 0) {
177 doLogInfo(workerID, "Waiting for answer");
178 // asking for primary generation
179
180 auto code = requestchannel.Receive(reply);
181 if (code > 0) {
182 doLogInfo(workerID, "Primary chunk received");
183 auto rawmessage = std::move(reply.At(0));
184 auto header = *(o2::PrimaryChunkAnswer*)(rawmessage->GetData());
185 if (!header.payload_attached) {
186 doLogInfo(workerID, "No payload; Server in stage " + std::string(PrimStateToString[(int)header.serverstate]));
187 // if no payload attached we inspect the server state, to see what to do
189 sleep(1); // back-off and retry
190 return true;
191 }
192 // we need to decide what to do when the server is idle ---> if this happens immediately after a new batch request it means that the server might just lag a bit behind
193 return false;
194 } else {
195 auto payload = std::move(reply.At(1));
196 // wrap incoming bytes as a TMessageWrapper which offers "adoption" of a buffer
197 auto message = new TMessageWrapper(payload->GetData(), payload->GetSize());
198 auto chunk = static_cast<o2::data::PrimaryChunk*>(message->ReadObjectAny(message->GetClass()));
199
200 bool goon = true;
201 // no particles and eventID == -1 --> indication for no more work
202 if (chunk->mParticles.size() == 0 && chunk->mSubEventInfo.eventID == -1) {
203 doLogInfo(workerID, "No particles in reply : quitting kernel");
204 goon = false;
205 }
206
207 if (goon) {
208
209 auto info = chunk->mSubEventInfo;
210 LOG(info) << workerStr() << " Processing " << chunk->mParticles.size() << " primary particles "
211 << "for event " << info.eventID << "/" << info.maxEvents << " "
212 << "part " << info.part << "/" << info.nparts;
213
214 if (eventselection == nullptr || (focus_on_event == info.eventID && focus_on_part == info.part)) {
215 mVMCApp->setPrimaries(chunk->mParticles);
216 } else {
217 // nothing to transport here
218 mVMCApp->setPrimaries(std::vector<TParticle>{});
219 LOG(info) << workerStr() << " This chunk will be skipped";
220 }
221
222 mVMCApp->setSubEventInfo(&info);
223
224 if (reproducibleSim) {
225 LOG(info) << workerStr() << " Setting seed for this sub-event to " << chunk->mSubEventInfo.seed;
226 gRandom->SetSeed(chunk->mSubEventInfo.seed);
228 }
229
230 // Process one event
231 auto& conf = o2::conf::SimConfig::Instance();
232 if (strcmp(conf.getMCEngine().c_str(), "TGeant4") == 0 || strcmp(conf.getMCEngine().c_str(), "O2TrivialMCEngine") == 0) {
233 // this is preferred and necessary for Geant4
234 // since repeated "ProcessRun" might have significant overheads
235 mVMC->ProcessEvent();
236 } else {
237 // for Geant3 calling ProcessEvent is not enough
238 // as some hooks are not called
239 mVMC->ProcessRun(1);
240 }
241
242 FairSystemInfo sysinfo;
243 LOG(info) << workerStr() << " TIME-STAMP " << mTimer.RealTime() << "\t";
244 mTimer.Continue();
245 LOG(info) << workerStr() << " MEM-STAMP " << sysinfo.GetCurrentMemory() / (1024. * 1024) << " "
246 << sysinfo.GetMaxMemory() << " MB\n";
247 }
248 delete message;
249 delete chunk;
250 }
251 } else {
252 LOG(info) << workerStr() << " No primary answer received from server (within timeout). Return code " << code;
253 }
254 } else {
255 LOG(info) << workerStr() << " Requesting work from server not possible. Return code " << sendcode;
256 return false;
257 }
258 return true;
259}
260
262{
263 return Kernel(-1, GetChannels().at("primary-get").at(0), GetChannels().at("simdata").at(0));
264}
265
266void O2SimDevice::PostRun() { LOG(info) << "Shutting down "; }
267
268} // namespace devices
269} // namespace o2
Definition of the Stack class.
std::vector< std::string > header
int32_t i
GPUTPCCFCheckPadBaseline Kernel
SurfaceTrackState state
void doLogInfo(int workerID, std::string const &message)
uint16_t pos
Definition RawData.h:3
std::vector< std::string > splitString(const std::string &src, char delim, bool trim=false)
A TMessage reading from a buffer it does not own.
static VMCSeederService const & instance()
static SimConfig & Instance()
Definition SimConfig.h:112
static bool initSim(fair::mq::Channel &channel, std::unique_ptr< FairRunSim > &simptr)
bool Kernel(int workerID, fair::mq::Channel &requestchannel, fair::mq::Channel &dataoutchannel, fair::mq::Channel *statuschannel=nullptr)
~O2SimDevice() final
Default destructor.
bool isWorkAvailable(fair::mq::Channel &statuschannel, int workerID=-1)
void InitTask() final
Overloads the InitTask() method of fair::mq::Device.
bool ConditionalRun() final
Overloads the ConditionalRun() method of fair::mq::Device.
void setSimDataChannel(fair::mq::Channel *channel)
void setPrimaries(std::vector< TParticle > const &p)
void setSubEventInfo(o2::data::SubEventInfo *i)
static ShmManager & Instance()
Definition ShmManager.h:61
GLuint GLsizei const GLchar * message
Definition glcorearb.h:2517
GLuint counter
Definition glcorearb.h:3987
a couple of static helper functions to create timestamp values for CCDB queries or override obsolete ...
bool querySimConfig(fair::mq::Channel &channel)
constexpr const char * PrimStateToString[5]
LOG(info)<< "Compressed in "<< sw.CpuTime()<< " s"
const std::string str