OPAL (Object Oriented Parallel Accelerator Library) 2024.2
OPAL
Pilot.h
Go to the documentation of this file.
1//
2// Class Pilot
3// The Optimization Pilot (Master): Coordinates requests by optimizer
4// to workers and reports results back on a given communicator.
5//
6// Every worker thread notifies the master here if idle or not. When
7// available the master dispatches one of the pending simulations to the
8// worker who will run the specified simulation and report results back to
9// the master. The Optimizer class will poll the scheduler to check if some
10// (or all) results are available and continue to optimize and request new
11// simulation results.
12//
13// @see Worker
14// @see Optimizer
15//
16// @tparam Opt_t type of the optimizer
17// @tparam Sim_t type of the simulation
18// @tparam SolPropagationGraph_t strategy to distribute solution between
19// master islands
20// @tparam Comm_t comm splitter strategy
21//
22// Copyright (c) 2010 - 2013, Yves Ineichen, ETH Zürich
23// All rights reserved
24//
25// Implemented as part of the PhD thesis
26// "Toward massively parallel multi-objective optimization with application to
27// particle accelerators" (https://doi.org/10.3929/ethz-a-009792359)
28//
29// This file is part of OPAL.
30//
31// OPAL is free software: you can redistribute it and/or modify
32// it under the terms of the GNU General Public License as published by
33// the Free Software Foundation, either version 3 of the License, or
34// (at your option) any later version.
35//
36// You should have received a copy of the GNU General Public License
37// along with OPAL. If not, see <https://www.gnu.org/licenses/>.
38//
39#ifndef __PILOT_H__
40#define __PILOT_H__
41
42#include <mpi.h>
43#include <iostream>
44#include <string>
45#include <unistd.h>
46
47#include "Comm/MasterNode.h"
48#include "Comm/CommSplitter.h"
49
50#include "Util/AnsiColors.h"
51#include "Util/Types.h"
52#include "Util/CmdArguments.h"
54
55#include "Pilot/Poller.h"
56#include "Pilot/Worker.h"
57#include "Optimizer/Optimizer.h"
58
59#include "Util/Trace/Trace.h"
60#include "Util/Trace/FileSink.h"
62
64
65
94template <
95 class Opt_t
96 , class Sim_t
97 , class SolPropagationGraph_t
98 , class Comm_t
99>
100class Pilot : protected Poller {
101
102public:
103
104 // constructor only for Pilot classes inherited from this class
105 // they have their own setup function
106 Pilot(CmdArguments_t args, std::shared_ptr<Comm_t> comm,
107 const DVarContainer_t &dvar)
108 : Poller(comm->mpiComm())
109 , comm_(comm)
110 , cmd_args_(args)
111 , dvars_(dvar)
112 {
113 // do nothing
114 }
115
116 Pilot(CmdArguments_t args, std::shared_ptr<Comm_t> comm,
117 functionDictionary_t known_expr_funcs)
118 : Poller(comm->mpiComm())
119 , comm_(comm)
120 , cmd_args_(args)
121 {
122 setup(known_expr_funcs);
123 }
124
125 Pilot(CmdArguments_t args, std::shared_ptr<Comm_t> comm,
126 functionDictionary_t known_expr_funcs,
127 const DVarContainer_t &dvar,
128 const Expressions::Named_t &obj,
129 const Expressions::Named_t &cons,
130 std::vector<double> hypervolRef = {},
131 bool isOptimizerRun = true,
132 const std::map<std::string, std::string> &userVariables = {})
133 : Poller(comm->mpiComm())
134 , comm_(comm)
135 , cmd_args_(args)
136 , objectives_(obj)
137 , constraints_(cons)
138 , dvars_(dvar)
139 , hypervolRef_(hypervolRef)
140 {
141 if (isOptimizerRun)
142 setup(known_expr_funcs, userVariables);
143 }
144
145 virtual ~Pilot()
146 {
147 for (auto itr = objectives_.begin(); itr != objectives_.end(); ++ itr)
148 delete itr->second;
149
150 for (auto itr = constraints_.begin(); itr != constraints_.end(); ++ itr)
151 delete itr->second;
152 }
153
154
155protected:
156
158 MPI_Comm worker_comm_;
160 MPI_Comm opt_comm_;
163
164 std::shared_ptr<Comm_t> comm_;
166
170
172
173 typedef MasterNode< typename Opt_t::SolutionState_t,
174 SolPropagationGraph_t > MasterNode_t;
175 std::unique_ptr< MasterNode_t > master_node_;
176
178 std::string input_file_;
179
183
187 std::vector<double> hypervolRef_;
188
189
190 // keep track of state of all workers
191 std::vector<bool> is_worker_idle_;
192
194 typedef std::map<size_t, std::pair<Param_t, reqVarContainer_t> > Jobs_t;
195 typedef Jobs_t::iterator JobIter_t;
198
199 //DEBUG
200 std::unique_ptr<Trace> job_trace_;
201
202private:
203 void setup(functionDictionary_t known_expr_funcs,
204 const std::map<std::string, std::string> &userVariables) {
205 global_rank_ = comm_->globalRank();
206
207 if(global_rank_ == 0) {
208 std::cout << AnsiColors::BoldMagenta;
209 std::cout << " _ _ _ _ " << std::endl;
210 std::cout << " | | (_) | | | " << std::endl;
211 std::cout << " ___ _ __ | |_ ______ _ __ _| | ___ | |_ " << std::endl;
212 std::cout << " / _ \\| '_ \\| __|______| '_ \\| | |/ _ \\| __|" << std::endl;
213 std::cout << "| (_) | |_) | |_ | |_) | | | (_) | |_ " << std::endl;
214 std::cout << " \\___/| .__/ \\__| | .__/|_|_|\\___/ \\__|" << std::endl;
215 std::cout << " | | | | " << std::endl;
216 std::cout << " |_| |_| " << std::endl;
217 // ADA std::cout << "☷ Version: \t" << PACKAGE_VERSION << std::endl;
218 //std::cout << "☷ Git: \t\t" << GIT_VERSION << std::endl;
219 //std::cout << "☷ Build Date: \t" << BUILD_DATE << std::endl;
220 std::cout << AnsiColors::Reset;
221 std::cout << std::endl;
222 }
223
224 MPI_Barrier(MPI_COMM_WORLD);
225 parseInputFile(known_expr_funcs, true);
226
227 // here the control flow starts to diverge
228 if ( comm_->isOptimizer() ) { startOptimizer(); }
229 else if ( comm_->isWorker() ) { startWorker(userVariables); }
230 else if ( comm_->isPilot() ) { startPilot(); }
231 }
232
233protected:
234
235 void parseInputFile(functionDictionary_t /*known_expr_funcs*/, bool isOptimizationRun) {
236
237 try {
238 input_file_ = cmd_args_->getArg<std::string>("inputfile", true);
239 } catch (OptPilotException &e) {
240 std::cout << "Could not find 'inputfile' in arguments.. Aborting."
241 << std::endl;
242 MPI_Abort(comm_m, -101);
243 }
244
245 if((isOptimizationRun && objectives_.size() == 0) || dvars_.size() == 0) {
246 throw OptPilotException("Pilot::Pilot()",
247 "No objectives or dvars specified");
248 }
249
250 if(global_rank_ == 0) {
251 std::ostringstream os;
253 os << " ✔ " << objectives_.size()
254 << " objectives" << std::endl;
255 if (isOptimizationRun) {
256 os << " ✔ " << constraints_.size()
257 << " constraints" << std::endl;
258 }
259 os << " ✔ " << dvars_.size()
260 << " dvars" << std::endl;
261 os << AnsiColors::Reset;
262 os << std::endl;
263 std::cout << os.str() << std::flush;
264 }
265
266 MPI_Barrier(MPI_COMM_WORLD);
267 }
268
269 virtual
271
272 std::ostringstream os;
273 os << AnsiColors::BoldMagenta << " " << global_rank_ << " (PID: " << getpid() << ") ▶ Opt"
274 << AnsiColors::Reset << std::endl;
275 std::cout << os.str() << std::flush;
276
277 const std::unique_ptr<Opt_t> opt(
278 new Opt_t(objectives_, constraints_, dvars_, objectives_.size(),
279 comm_->getBundle(), cmd_args_, hypervolRef_, comm_->getNrWorkerGroups()));
280 opt->initialize();
281
282 std::cout << "Stop Opt.." << std::endl;
283 }
284
285 virtual
286 void startWorker(const std::map<std::string, std::string> &userVariables) {
287
288 std::ostringstream os;
289 os << AnsiColors::BoldMagenta << " " << global_rank_ << " (PID: " << getpid() << ") ▶ Worker"
290 << AnsiColors::Reset << std::endl;
291 std::cout << os.str() << std::flush;
292
293 size_t pos = input_file_.find_last_of("/");
294 std::string tmplfile = input_file_;
295 if(pos != std::string::npos)
296 tmplfile = input_file_.substr(pos+1);
297 pos = tmplfile.find(".");
298 std::string simName = tmplfile.substr(0,pos);
299
300 const std::unique_ptr< Worker<Sim_t> > w(
302 comm_->getBundle(), cmd_args_, userVariables));
303
304 std::cout << "Stop Worker.." << std::endl;
305 }
306
307 virtual
308 void startPilot() {
309
310 std::ostringstream os;
311 os << AnsiColors::BoldMagenta << " " << global_rank_ << " (PID: " << getpid() << ") ▶ Pilot"
312 << AnsiColors::Reset << std::endl;
313 std::cout << os.str() << std::flush;
314
315 // Traces
316 std::ostringstream trace_filename;
317 trace_filename << "pilot.trace." << comm_->getBundle().island_id;
318 job_trace_.reset(new Trace("Optimizer Job Trace"));
319 job_trace_->registerComponent( "sink",
320 std::shared_ptr<TraceComponent>(new FileSink(trace_filename.str())));
321
322 worker_comm_ = comm_->getBundle().worker;
323 opt_comm_ = comm_->getBundle().opt;
324 coworker_comm_ = comm_->getBundle().world;
325
327 MPI_Comm_rank(worker_comm_, &my_rank_in_worker_comm_);
329 MPI_Comm_rank(opt_comm_, &my_rank_in_opt_comm_);
330
332 MPI_Comm_size(worker_comm_, &total_available_workers_);
335
336 // setup master network
337 num_coworkers_ = 0;
338 MPI_Comm_size(coworker_comm_, &num_coworkers_);
339 if(num_coworkers_ > 1) {
340 //FIXME: proper upper bound for window size
341 int alpha = cmd_args_->getArg<int>("initialPopulation", false);
342 int opt_size = objectives_.size() + constraints_.size();
343 int overhead = 10;
344 size_t upperbound_buffer_size =
345 sizeof(double) * alpha * (1 + opt_size) * 1000
346 + overhead;
347 master_node_.reset(
348 new MasterNode< typename Opt_t::SolutionState_t,
349 SolPropagationGraph_t >(
350 coworker_comm_, upperbound_buffer_size, objectives_.size(),
351 comm_->getBundle().island_id));
352 }
353
354 has_opt_converged_ = false;
355 continue_polling_ = true;
356 run();
357
358 std::cout << "Stop Pilot.." << std::endl;
359 }
360
361 virtual
363 {}
364
365 virtual
366 void prePoll()
367 {}
368
369 virtual
370 void onStop()
371 {}
372
373 virtual
374 void postPoll() {
375 // terminating all workers is tricky since we do not know their state.
376 // All workers are notified (to terminate) when opt has converged and
377 // all workers are idle.
378 bool all_worker_idle = true;
379
380 // in the case where new requests became available after worker
381 // delivered last results (and switched to idle state).
382 for(int i = 0; i < total_available_workers_; i++) {
383
384 if(i == my_rank_in_worker_comm_) continue;
385
386 all_worker_idle = all_worker_idle && is_worker_idle_[i];
387
388 if(is_worker_idle_[i] && !request_queue_.empty())
390 }
391
392 // when all workers have been notified we can stop polling
393 if(all_worker_idle && has_opt_converged_) {
394 continue_polling_ = false;
395 int dummy = 0;
396 for(int worker = 0; worker < total_available_workers_; worker++) {
397 MPI_Request req;
398 MPI_Isend(&dummy, 1, MPI_INT, worker,
400 }
401 }
402 }
403
404
405 virtual
406 void sendNewJobToWorker(int worker) {
407
408 // no new jobs once our opt has converged
409 if(has_opt_converged_) return;
410
411 JobIter_t job = request_queue_.begin();
412 size_t jid = job->first;
413
414 Param_t job_params = job->second.first;
415 MPI_Send(&jid, 1, MPI_UNSIGNED_LONG, worker, MPI_WORK_JOBID_TAG, worker_comm_);
416 MPI_Send_params(job_params, worker, worker_comm_);
417
418 //reqVarContainer_t job_reqvars = job->second.second;
419 //MPI_Send_reqvars(job_reqvars, worker, worker_comm_);
420
421 running_job_list_.insert(std::pair<size_t,
422 std::pair<Param_t, reqVarContainer_t> >(job->first, job->second));
423 request_queue_.erase(jid);
424 is_worker_idle_[worker] = false;
425
426 std::ostringstream dump;
427 dump << "sent job with ID " << jid << " to worker " << worker
428 << std::endl;
429 job_trace_->log(dump);
430
431 }
432
433
434 virtual
435 bool onMessage(MPI_Status status, size_t recv_value){
436
437 MPITag_t tag = MPITag_t(status.MPI_TAG);
438 switch(tag) {
439
440 case WORKER_FINISHED_TAG: {
441
442 size_t job_id = recv_value;
443
444 size_t dummy = 1;
445 MPI_Send(&dummy, 1, MPI_UNSIGNED_LONG, status.MPI_SOURCE,
447
449 MPI_Recv_reqvars(res, status.MPI_SOURCE, worker_comm_);
450
451 running_job_list_.erase(job_id);
452 is_worker_idle_[status.MPI_SOURCE] = true;
453
454 std::ostringstream dump;
455 dump << "worker finished job with ID " << job_id << std::endl;
456 job_trace_->log(dump);
457
458
459 // optimizer already terminated, cannot accept new messages
460 if(has_opt_converged_) return true;
461
462 int opt_master_rank = comm_->getLeader();
463 MPI_Send(&job_id, 1, MPI_UNSIGNED_LONG, opt_master_rank,
465
466 MPI_Send_reqvars(res, opt_master_rank, opt_comm_);
467
468 // we keep worker busy _after_ results have been sent to optimizer
469 if(!request_queue_.empty())
470 sendNewJobToWorker(status.MPI_SOURCE);
471
472 return true;
473 }
474
475 case OPT_NEW_JOB_TAG: {
476
477 size_t job_id = recv_value;
478 int opt_master_rank = comm_->getLeader();
479
480 Param_t job_params;
481 MPI_Recv_params(job_params, (size_t)opt_master_rank, opt_comm_);
482
483 reqVarContainer_t reqVars;
484 //MPI_Recv_reqvars(reqVars, (size_t)opt_master_rank, job_size, opt_comm_);
485
486 std::pair<Param_t, reqVarContainer_t> job =
487 std::pair<Param_t, reqVarContainer_t>(job_params, reqVars);
488 request_queue_.insert(
489 std::pair<size_t, std::pair<Param_t, reqVarContainer_t> >(
490 job_id, job));
491
492 std::ostringstream dump;
493 dump << "new opt job with ID " << job_id << std::endl;
494 job_trace_->log(dump);
495
496 return true;
497 }
498
500
501 if(num_coworkers_ <= 1) return true;
502
503 std::ostringstream dump;
504 dump << "starting solution exchange.. " << status.MPI_SOURCE << std::endl;
505 job_trace_->log(dump);
506
507 // we start by storing or local solution state
508 size_t buffer_size = recv_value;
509 int opt_master_rank = status.MPI_SOURCE; //comm_->getLeader();
510
511 char *buffer = new char[buffer_size];
512 MPI_Recv(buffer, buffer_size, MPI_CHAR, opt_master_rank,
514 master_node_->store(buffer, buffer_size);
515 delete[] buffer;
516
517 dump.clear();
518 dump.str(std::string());
519 dump << "getting " << buffer_size << " bytes from OPT "
520 << opt_master_rank << std::endl;
521 job_trace_->log(dump);
522
523 // and then continue collecting all other solution states
524 std::ostringstream states;
525 master_node_->collect(states);
526 buffer_size = states.str().length();
527
528 dump.clear();
529 dump.str(std::string());
530 dump << "collected solution states of other PILOTS: "
531 << buffer_size << " bytes" << std::endl;
532 job_trace_->log(dump);
533
534 // send collected solution states to optimizer;
535 MPI_Send(&buffer_size, 1, MPI_UNSIGNED_LONG, opt_master_rank,
537
538 buffer = new char[buffer_size];
539 std::memcpy(buffer, states.str().c_str(), buffer_size);
540 MPI_Send(buffer, buffer_size, MPI_CHAR, opt_master_rank,
542
543 dump.clear();
544 dump.str(std::string());
545 dump << "sent set of new solutions to OPT" << std::endl;
546 job_trace_->log(dump);
547
548 delete[] buffer;
549
550 return true;
551 }
552
553 case OPT_CONVERGED_TAG: {
554 return stop();
555 }
556
558 is_worker_idle_[status.MPI_SOURCE] = true;
559 return true;
560 }
561
562 default: {
563 std::string msg = "(Pilot) Error: unexpected MPI_TAG: ";
564 msg += status.MPI_TAG;
565 throw OptPilotException("Pilot::onMessage", msg);
566 }
567 }
568 }
569
570 bool stop(bool isOpt = true) {
571
572 if(has_opt_converged_) return true;
573
574 has_opt_converged_ = true;
575 request_queue_.clear();
576 size_t dummy = 0;
577 MPI_Request req;
578 MPI_Isend(&dummy, 1, MPI_UNSIGNED_LONG, comm_->getLeader(), MPI_STOP_TAG, opt_comm_, &req);
579
580 if(! isOpt) return true;
581 if(num_coworkers_ <= 1) return true;
582
583 if(! cmd_args_->getArg<bool>("one-pilot-converge", false, false))
584 return true;
585
586 // propagate converged message to other pilots
587 // FIXME what happens if two island converge at the same time?
588 int my_rank = 0;
589 MPI_Comm_rank(coworker_comm_, &my_rank);
590 for(int i=0; i < num_coworkers_; i++) {
591 if(i == my_rank) continue;
592 MPI_Request req;
593 MPI_Isend(&dummy, 1, MPI_UNSIGNED_LONG, i, OPT_CONVERGED_TAG, coworker_comm_, &req);
594 }
595
596 return true;
597 }
598
599
600 // we overwrite run here to handle polling on two different communicators
601 //XXX: would be nice to give the poller interface an array of comms and
602 // listeners to be called..
603 void run() {
604
605 MPI_Request opt_request;
606 MPI_Request worker_request;
607 MPI_Status status;
608 int flag = 0;
609 size_t recv_value_worker = 0;
610 size_t recv_value_opt = 0;
611
612 setupPoll();
613
614 MPI_Irecv(&recv_value_opt, 1, MPI_UNSIGNED_LONG, MPI_ANY_SOURCE,
615 MPI_ANY_TAG, opt_comm_, &opt_request);
616 MPI_Irecv(&recv_value_worker, 1, MPI_UNSIGNED_LONG, MPI_ANY_SOURCE,
617 MPI_ANY_TAG, worker_comm_, &worker_request);
618
619 bool pending_opt_request = true;
620 bool pending_worker_request = true;
621 bool pending_pilot_request = false;
622
623 MPI_Request pilot_request;
624 size_t recv_value_pilot = 0;
625 if(cmd_args_->getArg<bool>("one-pilot-converge", false, false)) {
626 MPI_Irecv(&recv_value_pilot, 1, MPI_UNSIGNED_LONG, MPI_ANY_SOURCE,
627 MPI_ANY_TAG, coworker_comm_, &pilot_request);
628 pending_pilot_request = true;
629 }
630
631 while(continue_polling_) {
632
633 prePoll();
634
635 if(opt_request != MPI_REQUEST_NULL) {
636 MPI_Test(&opt_request, &flag, &status);
637 if(flag) {
638 pending_opt_request = false;
639 if(status.MPI_TAG == MPI_STOP_TAG) {
640 return;
641 } else {
642 if(onMessage(status, recv_value_opt)) {
643 MPI_Irecv(&recv_value_opt, 1, MPI_UNSIGNED_LONG,
644 MPI_ANY_SOURCE, MPI_ANY_TAG, opt_comm_,
645 &opt_request);
646 pending_opt_request = true;
647 } else
648 return;
649 }
650 }
651 }
652
653 if(worker_request != MPI_REQUEST_NULL) {
654 MPI_Test(&worker_request, &flag, &status);
655 if(flag) {
656 pending_worker_request = false;
657 if(status.MPI_TAG == MPI_STOP_TAG) {
658 return;
659 } else {
660 if(onMessage(status, recv_value_worker)) {
661 MPI_Irecv(&recv_value_worker, 1,
662 MPI_UNSIGNED_LONG, MPI_ANY_SOURCE, MPI_ANY_TAG,
663 worker_comm_, &worker_request);
664 pending_worker_request = true;
665 } else
666 return;
667 }
668 }
669 }
670
671 if(cmd_args_->getArg<bool>("one-pilot-converge", false, false)) {
672 if(pilot_request != MPI_REQUEST_NULL) {
673 MPI_Test(&pilot_request, &flag, &status);
674 if(flag) {
675 pending_pilot_request = false;
676 if(status.MPI_TAG == OPT_CONVERGED_TAG) {
677 stop(false);
678 } else {
679 MPI_Irecv(&recv_value_pilot, 1,
680 MPI_UNSIGNED_LONG, MPI_ANY_SOURCE, MPI_ANY_TAG,
681 coworker_comm_, &pilot_request);
682 pending_pilot_request = true;
683 }
684 }
685 }
686 }
687
688 postPoll();
689 }
690
691 if(pending_opt_request) MPI_Cancel( &opt_request );
692 if(pending_worker_request) MPI_Cancel( &worker_request );
693 if(pending_pilot_request) MPI_Cancel( &pilot_request );
694 }
695
696};
697
698#endif
std::map< std::string, client::function::type > functionDictionary_t
Definition Expression.h:51
std::map< std::string, DVar_t > DVarContainer_t
Definition Types.h:108
std::map< std::string, reqVarInfo_t > reqVarContainer_t
Definition Types.h:96
namedVariableCollection_t Param_t
Definition Types.h:52
void MPI_Send_reqvars(reqVarContainer_t reqvars, std::size_t pid, MPI_Comm comm)
void MPI_Send_params(Param_t params, std::size_t pid, MPI_Comm comm)
void MPI_Recv_reqvars(reqVarContainer_t &reqvars, std::size_t pid, MPI_Comm comm)
void MPI_Recv_params(Param_t &params, std::size_t pid, MPI_Comm comm)
std::shared_ptr< CmdArguments > CmdArguments_t
#define MPI_WORK_JOBID_TAG
unique id of the job
Definition MPIHelper.h:52
#define MPI_EXCHANGE_SOL_STATE_DATA_TAG
Definition MPIHelper.h:59
#define MPI_OPT_JOB_FINISHED_TAG
pilot tells optimizer that results are ready to collect
Definition MPIHelper.h:46
#define MPI_EXCHANGE_SOL_STATE_RES_SIZE_TAG
Definition MPIHelper.h:60
#define MPI_EXCHANGE_SOL_STATE_RES_TAG
Definition MPIHelper.h:61
MPITag_t
Definition MPIHelper.h:71
@ WORKER_FINISHED_TAG
Definition MPIHelper.h:72
@ OPT_CONVERGED_TAG
Definition MPIHelper.h:74
@ OPT_NEW_JOB_TAG
Definition MPIHelper.h:73
@ WORKER_STATUSUPDATE_TAG
Definition MPIHelper.h:75
@ EXCHANGE_SOL_STATE_TAG
Definition MPIHelper.h:77
#define MPI_WORKER_FINISHED_ACK_TAG
pilot notifies worker that he is ready to collect the results
Definition MPIHelper.h:37
#define MPI_STOP_TAG
global stop tag to exit poll loop (
Definition MPIHelper.h:64
std::map< std::string, Expressions::Expr_t * > Named_t
type of an expressions with a name
Definition Expression.h:68
constexpr char Reset[]
Definition AnsiColors.h:5
constexpr char BoldMagenta[]
Definition AnsiColors.h:7
Definition Pilot.h:100
bool has_opt_converged_
Definition Pilot.h:181
CmdArguments_t cmd_args_
Definition Pilot.h:165
int total_available_workers_
Definition Pilot.h:180
MPI_Comm coworker_comm_
MPI communicator used for messages between all pilots.
Definition Pilot.h:162
void setup(functionDictionary_t known_expr_funcs, const std::map< std::string, std::string > &userVariables)
Definition Pilot.h:203
std::string input_file_
input file for simulation with embedded optimization problem
Definition Pilot.h:178
std::map< size_t, std::pair< Param_t, reqVarContainer_t > > Jobs_t
keep track of requests and running jobs
Definition Pilot.h:194
virtual void setupPoll()
executed before starting polling loop
Definition Pilot.h:362
Jobs_t request_queue_
Definition Pilot.h:197
std::vector< bool > is_worker_idle_
Definition Pilot.h:191
virtual void startPilot()
Definition Pilot.h:308
int global_rank_
Definition Pilot.h:167
virtual void startWorker(const std::map< std::string, std::string > &userVariables)
Definition Pilot.h:286
virtual ~Pilot()
Definition Pilot.h:145
virtual void startOptimizer()
Definition Pilot.h:270
bool stop(bool isOpt=true)
Definition Pilot.h:570
virtual void prePoll()
executed before checking for new request
Definition Pilot.h:366
Jobs_t::iterator JobIter_t
Definition Pilot.h:195
int my_rank_in_worker_comm_
Definition Pilot.h:168
std::unique_ptr< Trace > job_trace_
Definition Pilot.h:200
Expressions::Named_t constraints_
constraints
Definition Pilot.h:185
virtual void onStop()
enable implementation to react to STOP tag
Definition Pilot.h:370
int my_rank_in_opt_comm_
Definition Pilot.h:169
virtual void sendNewJobToWorker(int worker)
Definition Pilot.h:406
std::unique_ptr< MasterNode_t > master_node_
Definition Pilot.h:175
MasterNode< typename Opt_t::SolutionState_t, SolPropagationGraph_t > MasterNode_t
Definition Pilot.h:174
Pilot(CmdArguments_t args, std::shared_ptr< Comm_t > comm, const DVarContainer_t &dvar)
Definition Pilot.h:106
int num_coworkers_
Definition Pilot.h:171
bool continue_polling_
Definition Pilot.h:182
std::vector< double > hypervolRef_
hypervolume reference point
Definition Pilot.h:187
MPI_Comm opt_comm_
MPI communicator used for messages to/from optimizer.
Definition Pilot.h:160
Expressions::Named_t objectives_
objectives
Definition Pilot.h:184
std::shared_ptr< Comm_t > comm_
Definition Pilot.h:164
MPI_Comm worker_comm_
MPI communicator used for messages to/from worker.
Definition Pilot.h:158
virtual void postPoll()
executed after handling (if any) new request
Definition Pilot.h:374
DVarContainer_t dvars_
design variables
Definition Pilot.h:186
void run()
Definition Pilot.h:603
virtual bool onMessage(MPI_Status status, size_t recv_value)
Definition Pilot.h:435
Pilot(CmdArguments_t args, std::shared_ptr< Comm_t > comm, functionDictionary_t known_expr_funcs)
Definition Pilot.h:116
Jobs_t running_job_list_
Definition Pilot.h:196
void parseInputFile(functionDictionary_t, bool isOptimizationRun)
Definition Pilot.h:235
Pilot(CmdArguments_t args, std::shared_ptr< Comm_t > comm, functionDictionary_t known_expr_funcs, const DVarContainer_t &dvar, const Expressions::Named_t &obj, const Expressions::Named_t &cons, std::vector< double > hypervolRef={}, bool isOptimizerRun=true, const std::map< std::string, std::string > &userVariables={})
Definition Pilot.h:125
MPI_Comm comm_m
communicator the poller listens to requests
Definition Poller.h:52
Definition Trace.h:31