From ed8d527b652441b7afcc2fd32331b5c27dced9b2 Mon Sep 17 00:00:00 2001 From: Eberhard Graether Date: Mon, 12 May 2014 17:35:55 +0200 Subject: [PATCH] utility: Added messaging system To use the messaging system the classes Message and MessageListener need to get derived. The Singelton MessageQueue is responsible for managing listeners and sending messages in a threadsafe way. review id = 9 fortune cookie message = You are on your track to starting something new. --- src/lib/CMakeLists.txt | 7 + src/lib/utility/messaging/Message.h | 30 +++ src/lib/utility/messaging/MessageBase.h | 16 ++ src/lib/utility/messaging/MessageListener.h | 28 +++ .../utility/messaging/MessageListenerBase.h | 27 +++ src/lib/utility/messaging/MessageQueue.cpp | 155 +++++++++++++ src/lib/utility/messaging/MessageQueue.h | 51 +++++ src/test/CMakeLists.txt | 1 + src/test/MessageQueueTestSuite.h | 203 ++++++++++++++++++ 9 files changed, 518 insertions(+) create mode 100644 src/lib/utility/messaging/Message.h create mode 100644 src/lib/utility/messaging/MessageBase.h create mode 100644 src/lib/utility/messaging/MessageListener.h create mode 100644 src/lib/utility/messaging/MessageListenerBase.h create mode 100644 src/lib/utility/messaging/MessageQueue.cpp create mode 100644 src/lib/utility/messaging/MessageQueue.h create mode 100644 src/test/MessageQueueTestSuite.h diff --git a/src/lib/CMakeLists.txt b/src/lib/CMakeLists.txt index 09a6eb50..d9b2dea5 100644 --- a/src/lib/CMakeLists.txt +++ b/src/lib/CMakeLists.txt @@ -103,6 +103,13 @@ add_files( utility/math/Vector2.h utility/math/Vector4.h + utility/messaging/Message.h + utility/messaging/MessageBase.h + utility/messaging/MessageListener.h + utility/messaging/MessageListenerBase.h + utility/messaging/MessageQueue.cpp + utility/messaging/MessageQueue.h + utility/ConfigManager.cpp utility/ConfigManager.h utility/Vector.h diff --git a/src/lib/utility/messaging/Message.h b/src/lib/utility/messaging/Message.h new file mode 100644 index 00000000..618087d0 --- /dev/null +++ b/src/lib/utility/messaging/Message.h @@ -0,0 +1,30 @@ +#ifndef MESSAGE_H +#define MESSAGE_H + +#include +#include + +#include "utility/messaging/MessageBase.h" +#include "utility/messaging/MessageQueue.h" + +template +class Message: public MessageBase +{ +public: + virtual ~Message() + { + } + + virtual std::string getType() const + { + return MessageType::getStaticType(); + } + + void dispatch() + { + std::shared_ptr message = std::make_shared(*dynamic_cast(this)); + MessageQueue::getInstance()->pushMessage(message); + } +}; + +#endif // MESSAGE_H diff --git a/src/lib/utility/messaging/MessageBase.h b/src/lib/utility/messaging/MessageBase.h new file mode 100644 index 00000000..a2ac09a2 --- /dev/null +++ b/src/lib/utility/messaging/MessageBase.h @@ -0,0 +1,16 @@ +#ifndef MESSAGE_BASE_H +#define MESSAGE_BASE_H + +#include + +class MessageBase +{ +public: + virtual ~MessageBase() + { + } + + virtual std::string getType() const = 0; +}; + +#endif // MESSAGE_BASE_H diff --git a/src/lib/utility/messaging/MessageListener.h b/src/lib/utility/messaging/MessageListener.h new file mode 100644 index 00000000..21cc8c6b --- /dev/null +++ b/src/lib/utility/messaging/MessageListener.h @@ -0,0 +1,28 @@ +#ifndef MESSAGE_LISTENER_H +#define MESSAGE_LISTENER_H + +#include + +#include "utility/messaging/MessageBase.h" +#include "utility/messaging/MessageListenerBase.h" +#include "utility/messaging/MessageQueue.h" + +template +class MessageListener: public MessageListenerBase +{ +public: + virtual std::string getType() const + { + return MessageType::getStaticType(); + } + + virtual void handleMessageBase(MessageBase* message) + { + handleMessage(dynamic_cast(message)); + } + +private: + virtual void handleMessage(MessageType* message) = 0; +}; + +#endif // MESSAGE_LISTENER_H diff --git a/src/lib/utility/messaging/MessageListenerBase.h b/src/lib/utility/messaging/MessageListenerBase.h new file mode 100644 index 00000000..b7535682 --- /dev/null +++ b/src/lib/utility/messaging/MessageListenerBase.h @@ -0,0 +1,27 @@ +#ifndef MESSAGE_LISTENER_BASE_H +#define MESSAGE_LISTENER_BASE_H + +#include + +#include "utility/messaging/MessageBase.h" +#include "utility/messaging/MessageQueue.h" + +class MessageListenerBase +{ +public: + MessageListenerBase() + { + MessageQueue::getInstance()->registerListener(this); + } + + virtual ~MessageListenerBase() + { + MessageQueue::getInstance()->unregisterListener(this); + } + + virtual std::string getType() const = 0; + + virtual void handleMessageBase(MessageBase*) = 0; +}; + +#endif // MESSAGE_LISTENER_BASE_H diff --git a/src/lib/utility/messaging/MessageQueue.cpp b/src/lib/utility/messaging/MessageQueue.cpp new file mode 100644 index 00000000..4542a1a0 --- /dev/null +++ b/src/lib/utility/messaging/MessageQueue.cpp @@ -0,0 +1,155 @@ +#include "utility/messaging/MessageQueue.h" + +#include + +#include "utility/logging/logging.h" +#include "utility/messaging/MessageBase.h" +#include "utility/messaging/MessageListenerBase.h" + +std::shared_ptr MessageQueue::getInstance() +{ + if (!s_instance) + { + s_instance = std::shared_ptr(new MessageQueue()); + } + + return s_instance; +} + +void MessageQueue::registerListener(MessageListenerBase* listener) +{ + std::lock_guard lock(m_listenersMutex); + m_listeners.push_back(listener); +} + +void MessageQueue::unregisterListener(MessageListenerBase* listener) +{ + std::lock_guard lock(m_listenersMutex); + for (size_t i = 0; i < m_listeners.size(); i++) + { + if (m_listeners[i] == listener) + { + m_listeners.erase(m_listeners.begin() + i); + + // m_currentListenerIndex and m_listenersLength need to be updated in case this happens while a message is + // handled. + if (i <= m_currentListenerIndex) + { + m_currentListenerIndex--; + } + + if (i < m_listenersLength) + { + m_listenersLength--; + } + + return; + } + } + + LOG_ERROR("Listener was not found"); +} + +void MessageQueue::pushMessage(std::shared_ptr message) +{ + std::lock_guard lock(m_backMessageBufferMutex); + m_backMessageBuffer->push(message); +} + +void MessageQueue::startMessageLoopThreaded() +{ + std::thread(&MessageQueue::startMessageLoop, this).detach(); +} + +void MessageQueue::startMessageLoop() +{ + { + std::lock_guard lock(m_loopMutex); + + if (m_loopIsRunning) + { + LOG_ERROR("Loop is already running"); + return; + } + + m_loopIsRunning = true; + } + + while (true) + { + { + std::lock_guard lock(m_loopMutex); + + if (!m_loopIsRunning) + { + return; + } + } + + processMessages(); + } +} + +void MessageQueue::stopMessageLoop() +{ + std::lock_guard lock(m_loopMutex); + + if (!m_loopIsRunning) + { + LOG_WARNING("Loop is not running"); + } + + m_loopIsRunning = false; +} + +bool MessageQueue::loopIsRunning() const +{ + std::lock_guard lock(m_loopMutex); + return m_loopIsRunning; +} + +std::shared_ptr MessageQueue::s_instance; + +MessageQueue::MessageQueue() + : m_currentListenerIndex(0) + , m_listenersLength(0) + , m_loopIsRunning(false) +{ + m_frontMessageBuffer = std::make_shared(); + m_backMessageBuffer = std::make_shared(); +} + +void MessageQueue::processMessages() +{ + { + std::lock_guard lock(m_backMessageBufferMutex); + m_backMessageBuffer.swap(m_frontMessageBuffer); + } + + while (m_frontMessageBuffer->size()) + { + std::shared_ptr message = m_frontMessageBuffer->front(); + m_frontMessageBuffer->pop(); + + std::lock_guard lock(m_listenersMutex); + + // m_listenersLength is saved, so that new listeners registered whithin message handling don't get the + // current message and the length can be reduced when a listener gets unregistered. + m_listenersLength = m_listeners.size(); + + // The currentListenerIndex holds the index of the current listener being handled, so it can be changed when a + // listener gets removed while message handling. + for (m_currentListenerIndex = 0; m_currentListenerIndex < m_listenersLength; m_currentListenerIndex++) + { + MessageListenerBase* listener = m_listeners[m_currentListenerIndex]; + + if (listener->getType() == message->getType()) + { + // The listenersMutex gets unlocked so changes to listeners are possible while message handling. + m_listenersMutex.unlock(); + listener->handleMessageBase(message.get()); + m_listenersMutex.lock(); + } + } + } +} diff --git a/src/lib/utility/messaging/MessageQueue.h b/src/lib/utility/messaging/MessageQueue.h new file mode 100644 index 00000000..593390c2 --- /dev/null +++ b/src/lib/utility/messaging/MessageQueue.h @@ -0,0 +1,51 @@ +#ifndef MESSAGE_QUEUE_H +#define MESSAGE_QUEUE_H + +#include +#include +#include + +class MessageBase; +class MessageListenerBase; + +class MessageQueue +{ +public: + static std::shared_ptr getInstance(); + + void registerListener(MessageListenerBase* listener); + void unregisterListener(MessageListenerBase* listener); + + void pushMessage(std::shared_ptr message); + + void startMessageLoopThreaded(); + void startMessageLoop(); + void stopMessageLoop(); + + bool loopIsRunning() const; + +private: + typedef std::queue > MessageBufferType; + + static std::shared_ptr s_instance; + + MessageQueue(); + MessageQueue(const MessageQueue&); + void operator=(const MessageQueue&); + + void processMessages(); + + std::shared_ptr m_frontMessageBuffer; + std::shared_ptr m_backMessageBuffer; + std::vector m_listeners; + + size_t m_currentListenerIndex; + size_t m_listenersLength; + bool m_loopIsRunning; + + std::mutex m_backMessageBufferMutex; + std::mutex m_listenersMutex; + mutable std::mutex m_loopMutex; +}; + +#endif // MESSAGE_QUEUE_H diff --git a/src/test/CMakeLists.txt b/src/test/CMakeLists.txt index 1ff8af7e..db11e981 100644 --- a/src/test/CMakeLists.txt +++ b/src/test/CMakeLists.txt @@ -4,6 +4,7 @@ add_files( ConfigManagerTestSuite.h CxxParserTestSuite.h LogManagerTestSuite.h + MessageQueueTestSuite.h TestSuiteFixture.cpp TestSuiteFixture.h Vector2TestSuite.h diff --git a/src/test/MessageQueueTestSuite.h b/src/test/MessageQueueTestSuite.h new file mode 100644 index 00000000..468e1405 --- /dev/null +++ b/src/test/MessageQueueTestSuite.h @@ -0,0 +1,203 @@ +#include + +#include + +#include "utility/messaging/Message.h" +#include "utility/messaging/MessageListener.h" +#include "utility/messaging/MessageQueue.h" + +class MessageQueueTestSuite: public CxxTest::TestSuite +{ +public: + void test_message_loop_starts_and_stops(void) + { + TS_ASSERT(!MessageQueue::getInstance()->loopIsRunning()); + + MessageQueue::getInstance()->startMessageLoopThreaded(); + + waitForThread(); + + TS_ASSERT(MessageQueue::getInstance()->loopIsRunning()); + + MessageQueue::getInstance()->stopMessageLoop(); + + waitForThread(); + + TS_ASSERT(!MessageQueue::getInstance()->loopIsRunning()); + } + + void test_registered_listener_receives_messages(void) + { + MessageQueue::getInstance()->startMessageLoopThreaded(); + + TestMessageListener listener; + Test2MessageListener listener2; + + TestMessage().dispatch(); + TestMessage().dispatch(); + TestMessage().dispatch(); + + waitForThread(); + + MessageQueue::getInstance()->stopMessageLoop(); + + TS_ASSERT_EQUALS(3, listener.m_messageCount); + TS_ASSERT_EQUALS(0, listener2.m_messageCount); + } + + void test_message_dispatching_within_message_handling(void) + { + MessageQueue::getInstance()->startMessageLoopThreaded(); + + TestMessageListener listener; + Test2MessageListener listener2; + + Test2Message().dispatch(); + + waitForThread(); + + MessageQueue::getInstance()->stopMessageLoop(); + + TS_ASSERT_EQUALS(1, listener.m_messageCount); + TS_ASSERT_EQUALS(1, listener2.m_messageCount); + } + + void test_listener_registration_within_message_handling(void) + { + MessageQueue::getInstance()->startMessageLoopThreaded(); + + Test3MessageListener listener; + + Test2Message().dispatch(); + TestMessage().dispatch(); + + waitForThread(); + + MessageQueue::getInstance()->stopMessageLoop(); + + TS_ASSERT(listener.m_listener); + if (listener.m_listener) + { + TS_ASSERT_EQUALS(1, listener.m_listener->m_messageCount); + } + } + + void test_listener_unregistration_within_message_handling(void) + { + MessageQueue::getInstance()->startMessageLoopThreaded(); + + Test4MessageListener listener; + + TestMessage().dispatch(); + + Test2Message().dispatch(); + + TestMessage().dispatch(); + TestMessage().dispatch(); + TestMessage().dispatch(); + + waitForThread(); + + MessageQueue::getInstance()->stopMessageLoop(); + + TS_ASSERT(listener.m_listener); + if (listener.m_listener) + { + TS_ASSERT_EQUALS(2, listener.m_listener->m_messageCount); + } + } + +private: + class TestMessage: public Message + { + public: + static const std::string getStaticType() + { + return "TestMessage"; + } + }; + + class Test2Message: public Message + { + public: + static const std::string getStaticType() + { + return "TestMessage2"; + } + }; + + class TestMessageListener: public MessageListener + { + public: + TestMessageListener() + : m_messageCount(0) + { + } + + int m_messageCount; + + private: + virtual void handleMessage(TestMessage* message) + { + m_messageCount++; + } + }; + + class Test2MessageListener: public MessageListener + { + public: + Test2MessageListener() + : m_messageCount(0) + { + } + + int m_messageCount; + + private: + virtual void handleMessage(Test2Message* message) + { + m_messageCount++; + TestMessage().dispatch(); + } + }; + + class Test3MessageListener: public MessageListener + { + public: + std::shared_ptr m_listener; + + private: + virtual void handleMessage(Test2Message* message) + { + m_listener = std::make_shared(); + } + }; + + class Test4MessageListener: + public MessageListener, + public MessageListener + { + public: + std::shared_ptr m_listener; + + private: + virtual void handleMessage(TestMessage* message) + { + if (!m_listener) + { + m_listener = std::make_shared(); + } + } + + virtual void handleMessage(Test2Message* message) + { + m_listener.reset(); + } + }; + + void waitForThread() const + { + static const int THREAD_WAIT_TIME_MS = 5; + std::this_thread::sleep_for(std::chrono::milliseconds(THREAD_WAIT_TIME_MS)); + } +};