122 using namespace fair::mq::hooks;
125 fair::mq::DeviceRunner runner{argc, argv};
127 runner.AddHook<SetCustomCmdLineOptions>([](DeviceRunner&
r) {
128 boost::program_options::options_description customOptions(
"Custom options");
130 r.fConfig.AddToCmdLineOptions(customOptions);
133 runner.AddHook<InstantiateDevice>([](DeviceRunner&
r) {
134 r.fDevice = std::unique_ptr<fair::mq::Device>{
getDevice(
r.fConfig)};
138 }
catch (std::exception& e) {
139 LOG(error) <<
"Unhandled exception reached the top of main: " << e.what()
140 <<
", application will now exit";
143 LOG(error) <<
"Non-exception instance being thrown. Please make sure you use std::runtime_exception() instead. "
144 <<
"Application will now exit.";
157KernelSetup initSim(std::string transport, std::string primaddress, std::string primstatusaddress, std::string mergeraddress,
int workerID)
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();
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();
168 auto datachannel =
new fair::mq::Channel{
"simdata",
"push", factory};
169 datachannel->Connect(mergeraddress);
170 datachannel->Validate();
177 return KernelSetup{sim, primchannel, datachannel, prim_status_channel, workerID};
229 static bool initialized =
false;
230 static fair::mq::Channel channel;
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));
241 std::unique_ptr<fair::mq::Message> reply(channel.NewMessage());
243 doLogInfo(workerID,
"Listening for master control input");
244 if (channel.Receive(reply) > 0) {
245 auto data = reply->GetData();
246 auto size = reply->GetSize();
248 std::string command(
reinterpret_cast<char const*
>(
data),
size);
249 doLogInfo(workerID,
"Received control message: " + command);
254 doLogInfo(workerID,
"Stop asked, shutting down");
259 doLogInfo(workerID,
"No control input received ");
264int main(
int argc,
char* argv[])
266 struct sigaction act;
267 memset(&act, 0,
sizeof act);
268 sigemptyset(&act.sa_mask);
270 act.sa_flags = SA_SIGINFO;
272 std::vector<int> handledsignals = {SIGTERM, SIGINT, SIGQUIT, SIGSEGV, SIGBUS, SIGFPE, SIGABRT};
274 for (
auto s : handledsignals) {
275 if (sigaction(s, &act,
nullptr)) {
276 LOG(error) <<
"Could not install signal handler for " << s;
283 fair::Logger::OnFatal([] {
throw fair::FatalException(
"Fatal error occured. Exiting without core dump..."); });
285 FairLogger::GetLogger()->SetLogVerbosityLevel(
"MEDIUM");
288 bpo::options_description desc{
"Options"};
291 (
"control",
"control type")
293 (
"config-key",
"config key")
294 (
"mq-config",bpo::value<std::string>(),
"path to FairMQ config")
295 (
"severity",
"log severity");
297 bpo::variables_map vm;
298 bpo::store(parse_command_line(argc, argv, desc), vm);
301 std::string FMQconfig;
302 if (vm.count(
"mq-config")) {
303 FMQconfig = vm[
"mq-config"].as<std::string>();
306 auto internalfork = getenv(
"ALICE_SIMFORKINTERNAL");
308 int driverPID = getppid();
311 if (FMQconfig.empty()) {
312 throw std::runtime_error(
"This should never be called without FairMQ 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;
324 std::string serveraddress;
325 std::string mergeraddress;
326 std::string serverstatus_address;
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();
343 sockets = (
channels[1])[
"sockets"].GetArray();
344 address = (sockets[0])[
"address"].GetString();
345 serverstatus_address =
address;
347 if (s ==
"hitmerger") {
348 auto channels = device[
"channels"].GetArray();
350 s = channel[
"name"].GetString();
351 if (s ==
"simdata") {
352 auto sockets = channel[
"sockets"].GetArray();
353 auto address = (sockets[0])[
"address"].GetString();
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.");
374 std::unique_ptr<FairRunSim> simrun;
376 if (!
initializeSim(
"zeromq", serverstatus_address, simrun)) {
377 LOG(error) <<
"Could not initialize simulation";
382 unsigned int nworkers = std::max(1u, std::thread::hardware_concurrency() / 2);
383 auto f = getenv(
"ALICE_NSIMWORKERS");
385 nworkers =
static_cast<unsigned int>(std::stoi(
f));
387 LOG(info) <<
"Running with " << nworkers <<
" sim workers ";
392 for (
auto i = 0u;
i < nworkers; ++
i) {
394 auto pid = (
i == nworkers - 1) ? 0 : fork();
400 auto collectAndPubThreadFunction = [driverPID, &pubchannel]() {
402 std::unique_ptr<fair::mq::Message>
msg(collectorchannel.NewMessage());
405 if (collectorchannel.Receive(
msg) > 0) {
408 std::string text(
reinterpret_cast<char const*
>(
data),
size);
414 if (
i == nworkers - 1) {
415 std::vector<std::thread> threads;
416 threads.push_back(std::thread(collectAndPubThreadFunction));
417 threads.back().detach();
427 auto kernelSetup =
initSim(
"zeromq", serveraddress, serverstatus_address, mergeraddress,
i);
429 std::stringstream worker;
430 worker <<
"WORKER" <<
i;
439 if (conf.asService()) {
440 LOG(info) <<
"IN SERVICE MODE WAITING";
447 LOG(info) <<
"FINISHING";
462 while ((cpid = wait(&status))) {
KernelSetup initSim(std::string transport, std::string primaddress, std::string primstatusaddress, std::string mergeraddress, int workerID)