Logo ROOT  
Reference Guide
 
Loading...
Searching...
No Matches
RooFit::MultiProcess::Messenger Class Reference

Manages the interprocess communication channels and wraps send and receive calls.

This class is used for all interprocess communication between the master, queue and worker processes. The communication runs over pipes built on socketpair(), which are created in the ProcessManager before forking, so that all processes inherit their ends of the connected channels; see Channel for the wire format.

Several channels connect the processes for different purposes:

  • The master and queue processes share a channel that is mainly used for sending tasks to the queue from master.
  • The queue process shares a channel with each worker process. These are used by the workers to obtain tasks from the queue.
  • The master shares a channel with each worker process. The master -> worker direction carries state updates (previously published over a ZeroMQ PUB-SUB socket) and the worker -> master direction carries back task results, which the master receives in 'JobManager::retrieve()'.
Parameters
process_managerProcessManager instance which manages the master, queue and worker processes that we want to set up communication for in this Messenger.

Definition at line 30 of file Messenger_decl.h.

Public Types

enum class  test_rcv_pipes { fromQonM , fromMonQ , fromWonQ , fromQonW }
 
enum class  test_snd_pipes { M2Q , Q2M , Q2W , W2Q }
 

Public Member Functions

 Messenger (ProcessManager &process_manager)
 
 ~Messenger ()
 
std::pair< Poller, std::size_t > create_queue_poller ()
 Helper function that creates a poller for Queue::loop()
 
std::pair< Poller, std::size_t > create_worker_poller ()
 Helper function that creates a poller for worker_loop()
 
template<typename T >
void publish_from_master_to_workers (T &&item)
 specialization that sends the final part of a message
 
template<typename T , typename T2 , typename... Ts>
void publish_from_master_to_workers (T &&item, T2 &&item2, Ts &&...items)
 specialization that sends the first parts of multipart messages
 
template<typename value_t >
value_t receive_from_master_on_queue ()
 
template<typename value_t >
value_t receive_from_master_on_worker (bool *more=nullptr)
 
template<typename value_t >
value_t receive_from_queue_on_master ()
 
template<typename value_t >
value_t receive_from_queue_on_worker ()
 
template<typename value_t >
value_t receive_from_worker_on_master (bool *more=nullptr)
 
template<typename value_t >
value_t receive_from_worker_on_queue (std::size_t this_worker_id)
 
void send_from_master_to_queue ()
 
template<typename T , typename... Ts>
void send_from_master_to_queue (T item, Ts... items)
 
void send_from_queue_to_master ()
 
template<typename T , typename... Ts>
void send_from_queue_to_master (T item, Ts... items)
 
void send_from_queue_to_worker (std::size_t this_worker_id)
 
template<typename T , typename... Ts>
void send_from_queue_to_worker (std::size_t this_worker_id, T item, Ts... items)
 
template<typename T >
void send_from_worker_to_master (T &&item)
 specialization that sends the final part of a message
 
template<typename T , typename T2 , typename... Ts>
void send_from_worker_to_master (T &&item, T2 &&item2, Ts &&...items)
 specialization that sends the first parts of multipart messages
 
void send_from_worker_to_queue ()
 
template<typename T , typename... Ts>
void send_from_worker_to_queue (T item, Ts... items)
 
void test_connections (const ProcessManager &process_manager)
 Test whether the channels between all processes are working.
 
void test_receive (X2X expected_ping_value, test_rcv_pipes rcv_pipe, std::size_t worker_id)
 
void test_send (X2X ping_value, test_snd_pipes snd_pipe, std::size_t worker_id)
 

Private Member Functions

void debug_print (std::string s)
 Function called from send and receive template functions in debug builds used to monitor the messages that are going to be sent or are received.
 
Channelselect_worker_channel_on_master ()
 On master: pick the worker channel to receive the next message from.
 
void update_worker_channel_on_master (Channel &channel, bool more)
 

Private Attributes

Channel mq_
 
std::vector< Channelmw_
 
Channelmw_current_source_ = nullptr
 
std::size_t mw_next_poll_position_ = 0
 
Poller mw_poller_
 
std::vector< Channelqw_
 
Channel this_worker_mw_
 
Channel this_worker_qw_
 

#include </github/home/ROOT-CI/src/roofit/multiprocess/res/RooFit/MultiProcess/Messenger_decl.h>

Member Enumeration Documentation

◆ test_rcv_pipes

Enumerator
fromQonM 
fromMonQ 
fromWonQ 
fromQonW 

Definition at line 44 of file Messenger_decl.h.

◆ test_snd_pipes

Enumerator
M2Q 
Q2M 
Q2W 
W2Q 

Definition at line 37 of file Messenger_decl.h.

Constructor & Destructor Documentation

◆ Messenger()

RooFit::MultiProcess::Messenger::Messenger ( ProcessManager & process_manager)
explicit

Definition at line 48 of file Messenger.cxx.

◆ ~Messenger()

RooFit::MultiProcess::Messenger::~Messenger ( )
default

Member Function Documentation

◆ create_queue_poller()

std::pair< Poller, std::size_t > RooFit::MultiProcess::Messenger::create_queue_poller ( )

Helper function that creates a poller for Queue::loop()

Definition at line 196 of file Messenger.cxx.

◆ create_worker_poller()

std::pair< Poller, std::size_t > RooFit::MultiProcess::Messenger::create_worker_poller ( )

Helper function that creates a poller for worker_loop()

Definition at line 207 of file Messenger.cxx.

◆ debug_print()

void RooFit::MultiProcess::Messenger::debug_print ( std::string s)
private

Function called from send and receive template functions in debug builds used to monitor the messages that are going to be sent or are received.

By defining this in the implementation file, compilation is a lot faster during debugging of Messenger or communication protocols.

Definition at line 307 of file Messenger.cxx.

◆ publish_from_master_to_workers() [1/2]

template<typename T >
void RooFit::MultiProcess::Messenger::publish_from_master_to_workers ( T && item)

specialization that sends the final part of a message

Definition at line 145 of file Messenger.h.

◆ publish_from_master_to_workers() [2/2]

template<typename T , typename T2 , typename... Ts>
void RooFit::MultiProcess::Messenger::publish_from_master_to_workers ( T && item,
T2 && item2,
Ts &&... items )

specialization that sends the first parts of multipart messages

Definition at line 160 of file Messenger.h.

◆ receive_from_master_on_queue()

template<typename value_t >
value_t RooFit::MultiProcess::Messenger::receive_from_master_on_queue ( )

Definition at line 128 of file Messenger.h.

◆ receive_from_master_on_worker()

template<typename value_t >
value_t RooFit::MultiProcess::Messenger::receive_from_master_on_worker ( bool * more = nullptr)

Definition at line 175 of file Messenger.h.

◆ receive_from_queue_on_master()

template<typename value_t >
value_t RooFit::MultiProcess::Messenger::receive_from_queue_on_master ( )

Definition at line 101 of file Messenger.h.

◆ receive_from_queue_on_worker()

template<typename value_t >
value_t RooFit::MultiProcess::Messenger::receive_from_queue_on_worker ( )

Definition at line 72 of file Messenger.h.

◆ receive_from_worker_on_master()

template<typename value_t >
value_t RooFit::MultiProcess::Messenger::receive_from_worker_on_master ( bool * more = nullptr)

Definition at line 216 of file Messenger.h.

◆ receive_from_worker_on_queue()

template<typename value_t >
value_t RooFit::MultiProcess::Messenger::receive_from_worker_on_queue ( std::size_t this_worker_id)

Definition at line 45 of file Messenger.h.

◆ select_worker_channel_on_master()

Channel & RooFit::MultiProcess::Messenger::select_worker_channel_on_master ( )
private

On master: pick the worker channel to receive the next message from.

Continues an in-progress multipart message from the same worker; otherwise waits for any worker and picks one round-robin.

Definition at line 215 of file Messenger.cxx.

◆ send_from_master_to_queue() [1/2]

void RooFit::MultiProcess::Messenger::send_from_master_to_queue ( )

Definition at line 253 of file Messenger.cxx.

◆ send_from_master_to_queue() [2/2]

template<typename T , typename... Ts>
void RooFit::MultiProcess::Messenger::send_from_master_to_queue ( T item,
Ts... items )

Definition at line 115 of file Messenger.h.

◆ send_from_queue_to_master() [1/2]

void RooFit::MultiProcess::Messenger::send_from_queue_to_master ( )

Definition at line 251 of file Messenger.cxx.

◆ send_from_queue_to_master() [2/2]

template<typename T , typename... Ts>
void RooFit::MultiProcess::Messenger::send_from_queue_to_master ( T item,
Ts... items )

Definition at line 88 of file Messenger.h.

◆ send_from_queue_to_worker() [1/2]

void RooFit::MultiProcess::Messenger::send_from_queue_to_worker ( std::size_t this_worker_id)

Definition at line 247 of file Messenger.cxx.

◆ send_from_queue_to_worker() [2/2]

template<typename T , typename... Ts>
void RooFit::MultiProcess::Messenger::send_from_queue_to_worker ( std::size_t this_worker_id,
T item,
Ts... items )

Definition at line 59 of file Messenger.h.

◆ send_from_worker_to_master() [1/2]

template<typename T >
void RooFit::MultiProcess::Messenger::send_from_worker_to_master ( T && item)

specialization that sends the final part of a message

Definition at line 190 of file Messenger.h.

◆ send_from_worker_to_master() [2/2]

template<typename T , typename T2 , typename... Ts>
void RooFit::MultiProcess::Messenger::send_from_worker_to_master ( T && item,
T2 && item2,
Ts &&... items )

specialization that sends the first parts of multipart messages

Definition at line 203 of file Messenger.h.

◆ send_from_worker_to_queue() [1/2]

void RooFit::MultiProcess::Messenger::send_from_worker_to_queue ( )

Definition at line 245 of file Messenger.cxx.

◆ send_from_worker_to_queue() [2/2]

template<typename T , typename... Ts>
void RooFit::MultiProcess::Messenger::send_from_worker_to_queue ( T item,
Ts... items )

Definition at line 32 of file Messenger.h.

◆ test_connections()

void RooFit::MultiProcess::Messenger::test_connections ( const ProcessManager & process_manager)

Test whether the channels between all processes are working.

Parameters
process_managerProcessManager object used to instantiate this object. Used to identify which process we are running on and hence which channels need to be tested.

Definition at line 138 of file Messenger.cxx.

◆ test_receive()

void RooFit::MultiProcess::Messenger::test_receive ( X2X expected_ping_value,
test_rcv_pipes rcv_pipe,
std::size_t worker_id )

Definition at line 101 of file Messenger.cxx.

◆ test_send()

void RooFit::MultiProcess::Messenger::test_send ( X2X ping_value,
test_snd_pipes snd_pipe,
std::size_t worker_id )

Definition at line 79 of file Messenger.cxx.

◆ update_worker_channel_on_master()

void RooFit::MultiProcess::Messenger::update_worker_channel_on_master ( Channel & channel,
bool more )
private

Definition at line 238 of file Messenger.cxx.

Member Data Documentation

◆ mq_

Channel RooFit::MultiProcess::Messenger::mq_
private

Definition at line 111 of file Messenger_decl.h.

◆ mw_

std::vector<Channel> RooFit::MultiProcess::Messenger::mw_
private

Definition at line 119 of file Messenger_decl.h.

◆ mw_current_source_

Channel* RooFit::MultiProcess::Messenger::mw_current_source_ = nullptr
private

Definition at line 124 of file Messenger_decl.h.

◆ mw_next_poll_position_

std::size_t RooFit::MultiProcess::Messenger::mw_next_poll_position_ = 0
private

Definition at line 125 of file Messenger_decl.h.

◆ mw_poller_

Poller RooFit::MultiProcess::Messenger::mw_poller_
private

Definition at line 123 of file Messenger_decl.h.

◆ qw_

std::vector<Channel> RooFit::MultiProcess::Messenger::qw_
private

Definition at line 114 of file Messenger_decl.h.

◆ this_worker_mw_

Channel RooFit::MultiProcess::Messenger::this_worker_mw_
private

Definition at line 120 of file Messenger_decl.h.

◆ this_worker_qw_

Channel RooFit::MultiProcess::Messenger::this_worker_qw_
private

Definition at line 115 of file Messenger_decl.h.

  • roofit/multiprocess/res/RooFit/MultiProcess/Messenger_decl.h
  • roofit/multiprocess/res/RooFit/MultiProcess/Messenger.h
  • roofit/multiprocess/src/Messenger.cxx