X-Git-Url: https://bilbo.iut-bm.univ-fcomte.fr/and/gitweb/loba.git/blobdiff_plain/404a8d5b50296756e0896c1914750fa235720946..4fe47748b018bd1d159839fe9e2bb6685a895e85:/communicator.cpp?ds=inline diff --git a/communicator.cpp b/communicator.cpp index 753b04a..90cf990 100644 --- a/communicator.cpp +++ b/communicator.cpp @@ -34,6 +34,7 @@ communicator::communicator() receiver_process = MSG_process_create("receiver", communicator::receiver_wrapper, this, MSG_host_self()); + xbt_cond_wait(cond, mutex); // wait for the receiver to be ready xbt_mutex_release(mutex); } @@ -82,13 +83,26 @@ void communicator::send(const char* dest, message* msg) sent_comm.push_back(comm); } -bool communicator::recv(message*& msg, m_host_t& from, bool wait) +bool communicator::recv(message*& msg, m_host_t& from, double timeout) { - if (wait) { + if (timeout != 0) { + volatile double deadline = + timeout > 0 ? MSG_get_clock() + timeout : 0.0; xbt_mutex_acquire(mutex); - while (received.empty()) { + while (received.empty() && (!deadline || deadline > MSG_get_clock())) { + xbt_ex_t e; DEBUG0("waiting for a message to come"); - xbt_cond_wait(cond, mutex); + TRY { + if (deadline) + xbt_cond_timedwait(cond, mutex, deadline - MSG_get_clock()); + else + xbt_cond_wait(cond, mutex); + } + CATCH (e) { + if (e.category != timeout_error) + RETHROW; + xbt_ex_free(e); + } } xbt_mutex_release(mutex); } @@ -156,6 +170,11 @@ int communicator::receiver() { ctrl_comm = MSG_task_irecv(&ctrl_task, get_ctrl_mbox()); data_comm = MSG_task_irecv(&data_task, get_data_mbox()); + DEBUG0("receiver ready"); + xbt_mutex_acquire(mutex); + xbt_cond_signal(cond); // signal master that we are ready + xbt_mutex_release(mutex); + xbt_dynar_t comms = xbt_dynar_new(sizeof(msg_comm_t), NULL); while (ctrl_comm || data_comm) { @@ -169,7 +188,9 @@ int communicator::receiver() if (ctrl_comm && comm_test_n_destroy(ctrl_comm)) { if (strcmp(MSG_task_get_name(ctrl_task), "finalize")) { DEBUG0("received message from ctrl"); + xbt_mutex_acquire(mutex); received.push(ctrl_task); + xbt_mutex_release(mutex); ctrl_task = NULL; ctrl_comm = MSG_task_irecv(&ctrl_task, get_ctrl_mbox()); } else { @@ -183,7 +204,9 @@ int communicator::receiver() if (data_comm && comm_test_n_destroy(data_comm)) { if (strcmp(MSG_task_get_name(data_task), "finalize")) { DEBUG0("received message from data"); + xbt_mutex_acquire(mutex); received.push(data_task); + xbt_mutex_release(mutex); data_task = NULL; data_comm = MSG_task_irecv(&data_task, get_data_mbox()); } else { @@ -194,7 +217,8 @@ int communicator::receiver() } } xbt_mutex_acquire(mutex); - xbt_cond_signal(cond); + if (!received.empty()) + xbt_cond_signal(cond); xbt_mutex_release(mutex); } xbt_dynar_free(&comms);