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:
| process_manager | ProcessManager 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. | |
| Channel & | select_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< Channel > | mw_ |
| Channel * | mw_current_source_ = nullptr |
| std::size_t | mw_next_poll_position_ = 0 |
| Poller | mw_poller_ |
| std::vector< Channel > | qw_ |
| Channel | this_worker_mw_ |
| Channel | this_worker_qw_ |
#include </github/home/ROOT-CI/src/roofit/multiprocess/res/RooFit/MultiProcess/Messenger_decl.h>
| Enumerator | |
|---|---|
| fromQonM | |
| fromMonQ | |
| fromWonQ | |
| fromQonW | |
Definition at line 44 of file Messenger_decl.h.
| Enumerator | |
|---|---|
| M2Q | |
| Q2M | |
| Q2W | |
| W2Q | |
Definition at line 37 of file Messenger_decl.h.
|
explicit |
Definition at line 48 of file Messenger.cxx.
|
default |
| 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.
| 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.
|
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.
| 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.
| 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.
| value_t RooFit::MultiProcess::Messenger::receive_from_master_on_queue | ( | ) |
Definition at line 128 of file Messenger.h.
| value_t RooFit::MultiProcess::Messenger::receive_from_master_on_worker | ( | bool * | more = nullptr | ) |
Definition at line 175 of file Messenger.h.
| value_t RooFit::MultiProcess::Messenger::receive_from_queue_on_master | ( | ) |
Definition at line 101 of file Messenger.h.
| value_t RooFit::MultiProcess::Messenger::receive_from_queue_on_worker | ( | ) |
Definition at line 72 of file Messenger.h.
| value_t RooFit::MultiProcess::Messenger::receive_from_worker_on_master | ( | bool * | more = nullptr | ) |
Definition at line 216 of file Messenger.h.
| value_t RooFit::MultiProcess::Messenger::receive_from_worker_on_queue | ( | std::size_t | this_worker_id | ) |
Definition at line 45 of file Messenger.h.
|
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.
| void RooFit::MultiProcess::Messenger::send_from_master_to_queue | ( | ) |
Definition at line 253 of file Messenger.cxx.
| void RooFit::MultiProcess::Messenger::send_from_master_to_queue | ( | T | item, |
| Ts... | items ) |
Definition at line 115 of file Messenger.h.
| void RooFit::MultiProcess::Messenger::send_from_queue_to_master | ( | ) |
Definition at line 251 of file Messenger.cxx.
| void RooFit::MultiProcess::Messenger::send_from_queue_to_master | ( | T | item, |
| Ts... | items ) |
Definition at line 88 of file Messenger.h.
| void RooFit::MultiProcess::Messenger::send_from_queue_to_worker | ( | std::size_t | this_worker_id | ) |
Definition at line 247 of file Messenger.cxx.
| 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.
| 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.
| 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.
| void RooFit::MultiProcess::Messenger::send_from_worker_to_queue | ( | ) |
Definition at line 245 of file Messenger.cxx.
| void RooFit::MultiProcess::Messenger::send_from_worker_to_queue | ( | T | item, |
| Ts... | items ) |
Definition at line 32 of file Messenger.h.
| void RooFit::MultiProcess::Messenger::test_connections | ( | const ProcessManager & | process_manager | ) |
Test whether the channels between all processes are working.
| process_manager | ProcessManager 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.
| 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.
| 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.
|
private |
Definition at line 238 of file Messenger.cxx.
|
private |
Definition at line 111 of file Messenger_decl.h.
|
private |
Definition at line 119 of file Messenger_decl.h.
|
private |
Definition at line 124 of file Messenger_decl.h.
|
private |
Definition at line 125 of file Messenger_decl.h.
|
private |
Definition at line 123 of file Messenger_decl.h.
|
private |
Definition at line 114 of file Messenger_decl.h.
|
private |
Definition at line 120 of file Messenger_decl.h.
|
private |
Definition at line 115 of file Messenger_decl.h.