diff --git a/src/ebusd/baseloop.cpp b/src/ebusd/baseloop.cpp index 7e5f8a35..140054bb 100644 --- a/src/ebusd/baseloop.cpp +++ b/src/ebusd/baseloop.cpp @@ -85,10 +85,8 @@ void BaseLoop::start() // send result to client result += '\n'; - Connection* connection = message->getConnection(); - connection->addResult(NetMessage(result)); - - delete message; + message->setResult(result); + message->sendSignal(); // stop daemon if (strcasecmp(data.c_str(), "STOP") == 0) diff --git a/src/ebusd/connection.cpp b/src/ebusd/connection.cpp index f591f4cf..239bb0a1 100644 --- a/src/ebusd/connection.cpp +++ b/src/ebusd/connection.cpp @@ -26,12 +26,6 @@ extern Logger& L; int Connection::m_sum = 0; -void Connection::addResult(NetMessage message) -{ - NetMessage* tmp = new NetMessage(NetMessage(message)); - m_netQueueResult.add(tmp); -} - void* Connection::run() { m_running = true; @@ -80,22 +74,21 @@ void* Connection::run() // send data data[datalen] = '\0'; - m_netQueueData->add(new NetMessage(data, this)); + NetMessage message(data); + m_netQueue->add(&message); // wait for result L.log(net, debug, "[%05d] wait for result", getID()); - NetMessage* message = m_netQueueResult.remove(); + message.waitSignal(); L.log(net, debug, "[%05d] result added", getID()); - std::string result(message->getData()); + std::string result = message.getResult(); if (m_socket->isValid() == true) m_socket->send(result.c_str(), result.size()); else break; - delete message; - } } diff --git a/src/ebusd/connection.h b/src/ebusd/connection.h index bf847d32..e860bc8c 100644 --- a/src/ebusd/connection.h +++ b/src/ebusd/connection.h @@ -31,9 +31,7 @@ class Connection : public Thread public: Connection(TCPSocket* socket, WQueue* netQueue) - : m_socket(socket), m_netQueueData(netQueue), m_running(false) { m_sum++; m_id = m_sum;} - - void addResult(NetMessage message); + : m_socket(socket), m_netQueue(netQueue), m_running(false) { m_sum++; m_id = m_sum;} void* run(); void stop() const { m_notify.notify(); } @@ -43,8 +41,7 @@ public: private: TCPSocket* m_socket; - WQueue* m_netQueueData; - WQueue m_netQueueResult; + WQueue* m_netQueue; Notify m_notify; bool m_running; int m_id; diff --git a/src/ebusd/netmessage.h b/src/ebusd/netmessage.h index 7b285392..5c3449ae 100644 --- a/src/ebusd/netmessage.h +++ b/src/ebusd/netmessage.h @@ -35,16 +35,27 @@ public: /** * @brief constructs a new instance with message and source client address. * @param data from client. - * @param connection to return result to correct client. */ - NetMessage(const std::string data, Connection* connection=NULL) - : m_data(data), m_connection(connection) {} + NetMessage(const std::string data) : m_data(data) + { + pthread_mutex_init(&m_mutex, NULL); + pthread_cond_init(&m_cond, NULL); + } + + /** + * @brief destructor. + */ + ~NetMessage() + { + pthread_mutex_destroy(&m_mutex); + pthread_cond_destroy(&m_cond); + } /** * @brief copy constructor. * @param src message object for copy. */ - NetMessage(const NetMessage& src) : m_data(src.m_data), m_connection(src.m_connection) {} + NetMessage(const NetMessage& src) : m_data(src.m_data) {} /** * @brief get the data string. @@ -53,17 +64,52 @@ public: std::string getData() const { return m_data; } /** - * @brief original connection. - * @return pointer to connection. + * @brief get the result string. + * @return the result string. */ - Connection* getConnection() const { return m_connection; } + std::string getResult() const { return m_result; } + + /** + * @brief set the result string. + * @return the result string. + */ + void setResult(const std::string result) { m_result = result; } + + /** + * @brief wait on notification. + */ + void waitSignal() + { + pthread_mutex_lock(&m_mutex); + + while (m_result.size() == 0) + pthread_cond_wait(&m_cond, &m_mutex); + + pthread_mutex_unlock(&m_mutex); + } + + /** + * @brief send notification. + */ + void sendSignal() + { + pthread_mutex_lock(&m_mutex); + pthread_cond_signal(&m_cond); + pthread_mutex_unlock(&m_mutex); + } private: - /** the data/message string */ + /** the data string */ std::string m_data; - /** the source connection */ - Connection* m_connection; + /** the result string */ + std::string m_result; + + /** mutex variable for exclusive lock */ + pthread_mutex_t m_mutex; + + /** condition variable for exclusive lock */ + pthread_cond_t m_cond; };