Project
Loading...
Searching...
No Matches
O2SimDeviceRunner.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 "SimSetup/SimSetup.h"
16#include <SimConfig/SimConfig.h>
18#include <TVirtualMC.h>
19#include <thread>
20#include <fairmq/DeviceRunner.h>
21#include <boost/program_options.hpp>
22#include <memory>
23#include <string>
24#include <fairmq/Channel.h>
25#include <fairlogger/Logger.h>
26#include <fairmq/Parts.h>
27#include <fairmq/TransportFactory.h>
28#include <TStopwatch.h>
29#include <sys/wait.h>
30#include <pthread.h> // to set cpu affinity
31#include <cmath>
32#include <csignal>
33#include <unistd.h>
35
36#include "rapidjson/document.h"
37#include "rapidjson/stringbuffer.h"
38#include "rapidjson/filereadstream.h"
39namespace bpo = boost::program_options;
40
41std::vector<int> gChildProcesses; // global vector of child pids
44
45// a handler for error/termination signals
46void sigaction_handler(int signal, siginfo_t* signal_info, void*)
47{
48 auto pid = getpid();
49 LOG(info) << pid << " caught signal " << signal << " from source " << signal_info->si_pid;
50 auto groupid = getpgrp();
51 if (pid == gMasterProcess && signal_info->si_pid != gDriverProcess) {
52 // master worker forwards signal to whole worker process group
53 // (do this only if not coming from gDriverProcess since this uses killpg and already affected all children)
54 killpg(pid, signal);
55 } else {
56 if (signal_info->si_pid != gDriverProcess) {
57 // forward to master worker if coming internally
58 kill(groupid, signal);
59 }
60 }
61
62 if (signal_info->si_pid == gDriverProcess || signal == SIGTERM) {
63 // signal was sent from driver process --> not error
64 // or it was a standard SIGTERM
65
66 // shut down before waiting, so that the master worker finalises (e.g. writes its scoring dumps) before the driver's kill timer
68 // need to wait for potential children before exiting itself
69 // ... in order to have correct resource accounting
70 int status, cpid;
71 while ((cpid = wait(&status))) {
72 if (cpid == -1) {
73 break;
74 }
75 }
76 _exit(0);
77 }
78
79 // we treat internal signal interruption as an error
80 // because only ordinary termination is good in the context of the distributed system
81 _exit(128 + signal);
82}
83
84void addCustomOptions(bpo::options_description& options)
85{
86}
87
88void CustomCleanup(void* data, void* hint) { delete static_cast<std::string*>(hint); }
89
90// this will initialize the simulation setup
91// once before initializing the actual FairMQ device
92bool initializeSim(std::string transport, std::string address, std::unique_ptr<FairRunSim>& simptr)
93{
94 // This needs an already running PrimaryServer
95 auto factory = fair::mq::TransportFactory::CreateTransportFactory(transport);
96 auto channel = fair::mq::Channel{"o2sim-primserv-info", "req", factory};
97 channel.Connect(address);
98 channel.Validate();
99
100 return o2::devices::O2SimDevice::initSim(channel, simptr);
101}
102
104{
105 auto app = static_cast<o2::steer::O2MCApplication*>(TVirtualMCApplication::Instance());
106 auto vmc = TVirtualMC::GetMC();
107
108 if (app == nullptr) {
109 LOG(warning) << "no vmc application found at this stage";
110 }
111 return new o2::devices::O2SimDevice(app, vmc);
112}
113
114fair::mq::Device* getDevice(const fair::mq::ProgOptions& config)
115{
116 return getDevice();
117}
118
119int initAndRunDevice(int argc, char* argv[])
120{
121 using namespace fair::mq;
122 using namespace fair::mq::hooks;
123
124 try {
125 fair::mq::DeviceRunner runner{argc, argv};
126
127 runner.AddHook<SetCustomCmdLineOptions>([](DeviceRunner& r) {
128 boost::program_options::options_description customOptions("Custom options");
129 addCustomOptions(customOptions);
130 r.fConfig.AddToCmdLineOptions(customOptions);
131 });
132
133 runner.AddHook<InstantiateDevice>([](DeviceRunner& r) {
134 r.fDevice = std::unique_ptr<fair::mq::Device>{getDevice(r.fConfig)};
135 });
136
137 return runner.Run();
138 } catch (std::exception& e) {
139 LOG(error) << "Unhandled exception reached the top of main: " << e.what()
140 << ", application will now exit";
141 return 1;
142 } catch (...) {
143 LOG(error) << "Non-exception instance being thrown. Please make sure you use std::runtime_exception() instead. "
144 << "Application will now exit.";
145 return 1;
146 }
147}
148
151 fair::mq::Channel* primchannel = nullptr;
152 fair::mq::Channel* datachannel = nullptr;
153 fair::mq::Channel* primstatuschannel = nullptr;
154 int workerID = -1;
155};
156
157KernelSetup initSim(std::string transport, std::string primaddress, std::string primstatusaddress, std::string mergeraddress, int workerID)
158{
159 auto factory = fair::mq::TransportFactory::CreateTransportFactory(transport);
160 auto primchannel = new fair::mq::Channel{"primary-get", "req", factory};
161 primchannel->Connect(primaddress);
162 primchannel->Validate();
163
164 auto prim_status_channel = new fair::mq::Channel{"o2sim-primserv-info", "req", factory};
165 prim_status_channel->Connect(primstatusaddress);
166 prim_status_channel->Validate();
167
168 auto datachannel = new fair::mq::Channel{"simdata", "push", factory};
169 datachannel->Connect(mergeraddress);
170 datachannel->Validate();
171 // the channels are setup
172
173 // init the sim object
174 auto sim = getDevice();
175 sim->lateInit();
176
177 return KernelSetup{sim, primchannel, datachannel, prim_status_channel, workerID};
178}
179
181{
182 // the simplified runloop
183 while (setup.sim->Kernel(setup.workerID, *setup.primchannel, *setup.datachannel, setup.primstatuschannel)) {
184 }
185 doLogInfo(setup.workerID, "simulation is done");
187 return 0;
188}
189
190void pinToCPU(unsigned int cpuid)
191{
192 auto affinity = getenv("ALICE_CPUAFFINITY");
193 if (affinity) {
194 // MacOS does not support this API so we add a protection
195#ifndef __APPLE__
196
197 pthread_t thread;
198
199 thread = pthread_self();
200
201 cpu_set_t cpuset;
202 CPU_ZERO(&cpuset);
203 CPU_SET(cpuid, &cpuset);
204
205 auto s = pthread_setaffinity_np(thread, sizeof(cpu_set_t), &cpuset);
206 if (s != 0) {
207 LOG(warning) << "FAILED TO SET PTHREAD AFFINITY";
208 }
209
210 /* Check the actual affinity mask assigned to the thread */
211 s = pthread_getaffinity_np(thread, sizeof(cpu_set_t), &cpuset);
212 if (s != 0) {
213 LOG(warning) << "FAILED TO GET PTHREAD AFFINITY";
214 }
215
216 for (int j = 0; j < CPU_SETSIZE; j++) {
217 if (CPU_ISSET(j, &cpuset)) {
218 LOG(info) << "ENABLED CPU " << j;
219 }
220 }
221#else
222 LOG(warn) << "CPU AFFINITY NOT IMPLEMENTED ON APPLE";
223#endif
224 }
225}
226
227bool waitForControlInput(int workerID)
228{
229 static bool initialized = false;
230 static fair::mq::Channel channel;
231 if (!initialized) {
232 // we do the channel connect and initialization only once
233 // (reducing the chances that we might miss a control message from the master)
234 static auto factory = fair::mq::TransportFactory::CreateTransportFactory("zeromq");
235 channel = fair::mq::Channel{"o2sim-control", "sub", factory};
236 auto controlsocketname = getenv("ALICE_O2SIMCONTROL");
237 channel.Connect(std::string(controlsocketname));
238 channel.Validate();
239 initialized = true;
240 }
241 std::unique_ptr<fair::mq::Message> reply(channel.NewMessage());
242
243 doLogInfo(workerID, "Listening for master control input");
244 if (channel.Receive(reply) > 0) {
245 auto data = reply->GetData();
246 auto size = reply->GetSize();
247
248 std::string command(reinterpret_cast<char const*>(data), size);
249 doLogInfo(workerID, "Received control message: " + command);
250
252 o2::conf::parseSimReconfigFromString(command, reconfig);
253 if (reconfig.stop) {
254 doLogInfo(workerID, "Stop asked, shutting down");
255 return false;
256 }
257 doLogInfo(workerID, "Asked to process " + std::to_string(reconfig.nEvents) + std::string(" new events"));
258 } else {
259 doLogInfo(workerID, "No control input received ");
260 }
261 return true;
262}
263
264int main(int argc, char* argv[])
265{
266 struct sigaction act;
267 memset(&act, 0, sizeof act);
268 sigemptyset(&act.sa_mask);
269 act.sa_sigaction = &sigaction_handler;
270 act.sa_flags = SA_SIGINFO; // <--- enable sigaction
271
272 std::vector<int> handledsignals = {SIGTERM, SIGINT, SIGQUIT, SIGSEGV, SIGBUS, SIGFPE, SIGABRT}; // <--- may need to be completed
273 // remember that SIGKILL can't be handled
274 for (auto s : handledsignals) {
275 if (sigaction(s, &act, nullptr)) {
276 LOG(error) << "Could not install signal handler for " << s;
277 exit(EXIT_FAILURE);
278 }
279 }
280
281 // set the fatal callback for the logger to not do a core dump (since this might interfere with process shutdown sequence
282 // since it calls ROOT::TSystem and further child processes)
283 fair::Logger::OnFatal([] { throw fair::FatalException("Fatal error occured. Exiting without core dump..."); });
284 // initialy set logger verbosity to medium
285 FairLogger::GetLogger()->SetLogVerbosityLevel("MEDIUM");
286
287 // extract the path to FairMQ config
288 bpo::options_description desc{"Options"};
289 // clang-format off
290 desc.add_options()
291 ("control","control type")
292 ("id","ID")
293 ("config-key","config key")
294 ("mq-config",bpo::value<std::string>(),"path to FairMQ config")
295 ("severity","log severity");
296 // clang-format on
297 bpo::variables_map vm;
298 bpo::store(parse_command_line(argc, argv, desc), vm);
299 bpo::notify(vm);
300
301 std::string FMQconfig;
302 if (vm.count("mq-config")) {
303 FMQconfig = vm["mq-config"].as<std::string>();
304 }
305
306 auto internalfork = getenv("ALICE_SIMFORKINTERNAL");
307 if (internalfork) {
308 int driverPID = getppid();
309 auto pubchannel = o2::simpubsub::createPUBChannel(o2::simpubsub::getPublishAddress("o2sim-worker-notifications", driverPID));
310
311 if (FMQconfig.empty()) {
312 throw std::runtime_error("This should never be called without FairMQ config.");
313 }
314 // read the JSON config
315 FILE* fp = fopen(FMQconfig.c_str(), "r");
316 constexpr unsigned short usmax = std::numeric_limits<unsigned short>::max() - 1;
317 char readBuffer[usmax];
318 rapidjson::FileReadStream is(fp, readBuffer, sizeof(readBuffer));
319 rapidjson::Document d;
320 d.ParseStream(is);
321 fclose(fp);
322
323 // retrieve correct server and merger URLs
324 std::string serveraddress;
325 std::string mergeraddress;
326 std::string serverstatus_address;
327 std::string s;
328
329 auto& options = d["fairMQOptions"];
330 assert(options.IsObject());
331 for (auto option = options.MemberBegin(); option != options.MemberEnd(); ++option) {
332 s = option->name.GetString();
333 if (s == "devices") {
334 assert(option->value.IsArray());
335 auto devices = option->value.GetArray();
336 for (auto& device : devices) {
337 s = device["id"].GetString();
338 if (s == "primary-server") {
339 auto channels = device["channels"].GetArray();
340 auto sockets = (channels[0])["sockets"].GetArray();
341 auto address = (sockets[0])["address"].GetString();
342 serveraddress = address;
343 sockets = (channels[1])["sockets"].GetArray();
344 address = (sockets[0])["address"].GetString();
345 serverstatus_address = address;
346 }
347 if (s == "hitmerger") {
348 auto channels = device["channels"].GetArray();
349 for (auto& channel : channels) {
350 s = channel["name"].GetString();
351 if (s == "simdata") {
352 auto sockets = channel["sockets"].GetArray();
353 auto address = (sockets[0])["address"].GetString();
354 mergeraddress = address;
355 }
356 }
357 }
358 }
359 }
360 }
361
362 LOG(info) << "Parsed primary server address " << serveraddress;
363 LOG(info) << "Parsed primary server status address " << serverstatus_address;
364 LOG(info) << "Parsed merger address " << mergeraddress;
365 if (serveraddress.empty() || mergeraddress.empty()) {
366 throw std::runtime_error("Could not determine server or merger URLs.");
367 }
368 // This is a solution based on initializing the simulation once
369 // and then fork the process to share the simulation memory across
370 // many processes. Here we are not using fair::mq::Devices and just setup
371 // some channels manually and do our own runloop.
372
373 // we init the simulation first
374 std::unique_ptr<FairRunSim> simrun;
375 // TODO: take the addresses from somewhere else
376 if (!initializeSim("zeromq", serverstatus_address, simrun)) {
377 LOG(error) << "Could not initialize simulation";
378 return 1;
379 }
380
381 // should be factored out?
382 unsigned int nworkers = std::max(1u, std::thread::hardware_concurrency() / 2);
383 auto f = getenv("ALICE_NSIMWORKERS");
384 if (f) {
385 nworkers = static_cast<unsigned int>(std::stoi(f));
386 }
387 LOG(info) << "Running with " << nworkers << " sim workers ";
388
389 gMasterProcess = getpid();
390 gDriverProcess = getppid();
391 // then we fork and create a device in each fork
392 for (auto i = 0u; i < nworkers; ++i) {
393 // we use the current process as one of the workers as it has nothing else to do
394 auto pid = (i == nworkers - 1) ? 0 : fork();
395 if (pid == 0) {
396 // Each worker can publish its progress/state on a ZMQ channel.
397 // We actually use a push/pull mechanism to collect all messages in the
398 // master worker which can then publish using PUB/SUB.
399 // auto factory = fair::mq::TransportFactory::CreateTransportFactory("zeromq");
400 auto collectAndPubThreadFunction = [driverPID, &pubchannel]() {
401 auto collectorchannel = o2::simpubsub::createPUBChannel(o2::simpubsub::getPublishAddress("o2sim-workerinternal", driverPID), "pull");
402 std::unique_ptr<fair::mq::Message> msg(collectorchannel.NewMessage());
403
404 while (true) {
405 if (collectorchannel.Receive(msg) > 0) {
406 auto data = msg->GetData();
407 auto size = msg->GetSize();
408 std::string text(reinterpret_cast<char const*>(data), size);
409 // LOG(info) << "Collector message: " << text;
410 o2::simpubsub::publishMessage(pubchannel, text);
411 }
412 }
413 };
414 if (i == nworkers - 1) { // <---- extremely important to take non-forked version since ZMQ sockets do not behave well on fork
415 std::vector<std::thread> threads;
416 threads.push_back(std::thread(collectAndPubThreadFunction));
417 threads.back().detach();
418 }
419
420 // everyone else is getting a push socket for notifications
421 auto pushchannel = o2::simpubsub::createPUBChannel(o2::simpubsub::getPublishAddress("o2sim-workerinternal", driverPID), "push");
422
423 // we will try to pin each worker to a particular CPU
424 // this can be made configurable via environment variables??
425 pinToCPU(i);
426
427 auto kernelSetup = initSim("zeromq", serveraddress, serverstatus_address, mergeraddress, i);
428
429 std::stringstream worker;
430 worker << "WORKER" << i;
431 o2::simpubsub::publishMessage(pushchannel, o2::simpubsub::simStatusString(worker.str(), "STATUS", "SETUP COMPLETED"));
432
433 auto& conf = o2::conf::SimConfig::Instance();
434
435 bool more = true;
436 while (more) {
437 runSim(kernelSetup);
438
439 if (conf.asService()) {
440 LOG(info) << "IN SERVICE MODE WAITING";
441 o2::simpubsub::publishMessage(pushchannel, o2::simpubsub::simStatusString(worker.str(), "STATUS", "AWAITING INPUT"));
442 more = waitForControlInput(kernelSetup.workerID);
443 usleep(100); // --> why? (probably to give the server some chance to come to a "serving" state)
444 } else {
445 o2::simpubsub::publishMessage(pushchannel, o2::simpubsub::simStatusString(worker.str(), "STATUS", "TERMINATING"));
446
447 LOG(info) << "FINISHING";
448 more = false;
449 }
450 }
451 sleep(10); // ---> give some time for message to be delivered to merger (destructing too early might affect the ZQM buffers)
452 // The process will in any case be terminated by the main o2-sim driver.
453
454 // destruct setup (using _exit due to problems in ROOT shutdown (segmentation violations)
455 // Clearly at some moment, a more robust solution would be appreciated
456 _exit(0);
457 } else {
458 gChildProcesses.push_back(pid);
459 }
460 }
461 int status, cpid;
462 while ((cpid = wait(&status))) {
463 // LOG(info) << "normal wait " << cpid << " returned ";
464 if (cpid == -1) {
465 break;
466 }
467 }
468 _exit(0);
469 } else {
470 // This the solution where we setup an ordinary fair::mq::Device
471 // (each if which will setup its own simulation). Parallelism
472 // is achieved outside by instantiating multiple device processes.
473 _exit(initAndRunDevice(argc, argv));
474 }
475}
int32_t i
std::vector< int > gChildProcesses
void sigaction_handler(int signal, siginfo_t *signal_info, void *)
KernelSetup initSim(std::string transport, std::string primaddress, std::string primstatusaddress, std::string mergeraddress, int workerID)
int gMasterProcess
void pinToCPU(unsigned int cpuid)
void CustomCleanup(void *data, void *hint)
bool initializeSim(std::string transport, std::string address, std::unique_ptr< FairRunSim > &simptr)
bool waitForControlInput(int workerID)
int initAndRunDevice(int argc, char *argv[])
int runSim(KernelSetup setup)
o2::devices::O2SimDevice * getDevice()
void addCustomOptions(bpo::options_description &options)
int gDriverProcess
void doLogInfo(int workerID, std::string const &message)
uint32_t j
Definition RawData.h:0
uint16_t pid
Definition RawData.h:2
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)
GLuint GLuint64EXT address
Definition glcorearb.h:5846
GLsizeiptr size
Definition glcorearb.h:659
GLdouble f
Definition glcorearb.h:310
GLboolean * data
Definition glcorearb.h:298
GLboolean r
Definition glcorearb.h:1233
bool parseSimReconfigFromString(std::string const &argumentstring, SimReconfigData &config)
std::string simStatusString(std::string const &origin, std::string const &topic, std::string const &message)
std::string getPublishAddress(std::string const &base, int pid=getpid())
bool publishMessage(fair::mq::Channel &channel, std::string const &message)
fair::mq::Channel createPUBChannel(std::string const &address, std::string const &type="pub")
std::string to_string(gsl::span< T, Size > span)
Definition common.h:52
fair::mq::Channel * datachannel
o2::devices::O2SimDevice * sim
fair::mq::Channel * primstatuschannel
fair::mq::Channel * primchannel
static void shutdown()
Definition SimSetup.cxx:77
TODO: Make this a base class of SimConfigData?
Definition SimConfig.h:205
#define main
LOG(info)<< "Compressed in "<< sw.CpuTime()<< " s"
std::vector< ChannelData > channels
uint64_t const void const *restrict const msg
Definition x9.h:153