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.
This commit is contained in:
@@ -103,6 +103,13 @@ add_files(
|
|||||||
utility/math/Vector2.h
|
utility/math/Vector2.h
|
||||||
utility/math/Vector4.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.cpp
|
||||||
utility/ConfigManager.h
|
utility/ConfigManager.h
|
||||||
utility/Vector.h
|
utility/Vector.h
|
||||||
|
|||||||
@@ -0,0 +1,30 @@
|
|||||||
|
#ifndef MESSAGE_H
|
||||||
|
#define MESSAGE_H
|
||||||
|
|
||||||
|
#include <memory>
|
||||||
|
#include <string>
|
||||||
|
|
||||||
|
#include "utility/messaging/MessageBase.h"
|
||||||
|
#include "utility/messaging/MessageQueue.h"
|
||||||
|
|
||||||
|
template<typename MessageType>
|
||||||
|
class Message: public MessageBase
|
||||||
|
{
|
||||||
|
public:
|
||||||
|
virtual ~Message()
|
||||||
|
{
|
||||||
|
}
|
||||||
|
|
||||||
|
virtual std::string getType() const
|
||||||
|
{
|
||||||
|
return MessageType::getStaticType();
|
||||||
|
}
|
||||||
|
|
||||||
|
void dispatch()
|
||||||
|
{
|
||||||
|
std::shared_ptr<MessageBase> message = std::make_shared<MessageType>(*dynamic_cast<MessageType*>(this));
|
||||||
|
MessageQueue::getInstance()->pushMessage(message);
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
#endif // MESSAGE_H
|
||||||
@@ -0,0 +1,16 @@
|
|||||||
|
#ifndef MESSAGE_BASE_H
|
||||||
|
#define MESSAGE_BASE_H
|
||||||
|
|
||||||
|
#include <string>
|
||||||
|
|
||||||
|
class MessageBase
|
||||||
|
{
|
||||||
|
public:
|
||||||
|
virtual ~MessageBase()
|
||||||
|
{
|
||||||
|
}
|
||||||
|
|
||||||
|
virtual std::string getType() const = 0;
|
||||||
|
};
|
||||||
|
|
||||||
|
#endif // MESSAGE_BASE_H
|
||||||
@@ -0,0 +1,28 @@
|
|||||||
|
#ifndef MESSAGE_LISTENER_H
|
||||||
|
#define MESSAGE_LISTENER_H
|
||||||
|
|
||||||
|
#include <string>
|
||||||
|
|
||||||
|
#include "utility/messaging/MessageBase.h"
|
||||||
|
#include "utility/messaging/MessageListenerBase.h"
|
||||||
|
#include "utility/messaging/MessageQueue.h"
|
||||||
|
|
||||||
|
template<typename MessageType>
|
||||||
|
class MessageListener: public MessageListenerBase
|
||||||
|
{
|
||||||
|
public:
|
||||||
|
virtual std::string getType() const
|
||||||
|
{
|
||||||
|
return MessageType::getStaticType();
|
||||||
|
}
|
||||||
|
|
||||||
|
virtual void handleMessageBase(MessageBase* message)
|
||||||
|
{
|
||||||
|
handleMessage(dynamic_cast<MessageType*>(message));
|
||||||
|
}
|
||||||
|
|
||||||
|
private:
|
||||||
|
virtual void handleMessage(MessageType* message) = 0;
|
||||||
|
};
|
||||||
|
|
||||||
|
#endif // MESSAGE_LISTENER_H
|
||||||
@@ -0,0 +1,27 @@
|
|||||||
|
#ifndef MESSAGE_LISTENER_BASE_H
|
||||||
|
#define MESSAGE_LISTENER_BASE_H
|
||||||
|
|
||||||
|
#include <string>
|
||||||
|
|
||||||
|
#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
|
||||||
@@ -0,0 +1,155 @@
|
|||||||
|
#include "utility/messaging/MessageQueue.h"
|
||||||
|
|
||||||
|
#include <thread>
|
||||||
|
|
||||||
|
#include "utility/logging/logging.h"
|
||||||
|
#include "utility/messaging/MessageBase.h"
|
||||||
|
#include "utility/messaging/MessageListenerBase.h"
|
||||||
|
|
||||||
|
std::shared_ptr<MessageQueue> MessageQueue::getInstance()
|
||||||
|
{
|
||||||
|
if (!s_instance)
|
||||||
|
{
|
||||||
|
s_instance = std::shared_ptr<MessageQueue>(new MessageQueue());
|
||||||
|
}
|
||||||
|
|
||||||
|
return s_instance;
|
||||||
|
}
|
||||||
|
|
||||||
|
void MessageQueue::registerListener(MessageListenerBase* listener)
|
||||||
|
{
|
||||||
|
std::lock_guard<std::mutex> lock(m_listenersMutex);
|
||||||
|
m_listeners.push_back(listener);
|
||||||
|
}
|
||||||
|
|
||||||
|
void MessageQueue::unregisterListener(MessageListenerBase* listener)
|
||||||
|
{
|
||||||
|
std::lock_guard<std::mutex> 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<MessageBase> message)
|
||||||
|
{
|
||||||
|
std::lock_guard<std::mutex> lock(m_backMessageBufferMutex);
|
||||||
|
m_backMessageBuffer->push(message);
|
||||||
|
}
|
||||||
|
|
||||||
|
void MessageQueue::startMessageLoopThreaded()
|
||||||
|
{
|
||||||
|
std::thread(&MessageQueue::startMessageLoop, this).detach();
|
||||||
|
}
|
||||||
|
|
||||||
|
void MessageQueue::startMessageLoop()
|
||||||
|
{
|
||||||
|
{
|
||||||
|
std::lock_guard<std::mutex> lock(m_loopMutex);
|
||||||
|
|
||||||
|
if (m_loopIsRunning)
|
||||||
|
{
|
||||||
|
LOG_ERROR("Loop is already running");
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
m_loopIsRunning = true;
|
||||||
|
}
|
||||||
|
|
||||||
|
while (true)
|
||||||
|
{
|
||||||
|
{
|
||||||
|
std::lock_guard<std::mutex> lock(m_loopMutex);
|
||||||
|
|
||||||
|
if (!m_loopIsRunning)
|
||||||
|
{
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
processMessages();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
void MessageQueue::stopMessageLoop()
|
||||||
|
{
|
||||||
|
std::lock_guard<std::mutex> lock(m_loopMutex);
|
||||||
|
|
||||||
|
if (!m_loopIsRunning)
|
||||||
|
{
|
||||||
|
LOG_WARNING("Loop is not running");
|
||||||
|
}
|
||||||
|
|
||||||
|
m_loopIsRunning = false;
|
||||||
|
}
|
||||||
|
|
||||||
|
bool MessageQueue::loopIsRunning() const
|
||||||
|
{
|
||||||
|
std::lock_guard<std::mutex> lock(m_loopMutex);
|
||||||
|
return m_loopIsRunning;
|
||||||
|
}
|
||||||
|
|
||||||
|
std::shared_ptr<MessageQueue> MessageQueue::s_instance;
|
||||||
|
|
||||||
|
MessageQueue::MessageQueue()
|
||||||
|
: m_currentListenerIndex(0)
|
||||||
|
, m_listenersLength(0)
|
||||||
|
, m_loopIsRunning(false)
|
||||||
|
{
|
||||||
|
m_frontMessageBuffer = std::make_shared<MessageBufferType>();
|
||||||
|
m_backMessageBuffer = std::make_shared<MessageBufferType>();
|
||||||
|
}
|
||||||
|
|
||||||
|
void MessageQueue::processMessages()
|
||||||
|
{
|
||||||
|
{
|
||||||
|
std::lock_guard<std::mutex> lock(m_backMessageBufferMutex);
|
||||||
|
m_backMessageBuffer.swap(m_frontMessageBuffer);
|
||||||
|
}
|
||||||
|
|
||||||
|
while (m_frontMessageBuffer->size())
|
||||||
|
{
|
||||||
|
std::shared_ptr<MessageBase> message = m_frontMessageBuffer->front();
|
||||||
|
m_frontMessageBuffer->pop();
|
||||||
|
|
||||||
|
std::lock_guard<std::mutex> 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();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,51 @@
|
|||||||
|
#ifndef MESSAGE_QUEUE_H
|
||||||
|
#define MESSAGE_QUEUE_H
|
||||||
|
|
||||||
|
#include <memory>
|
||||||
|
#include <mutex>
|
||||||
|
#include <queue>
|
||||||
|
|
||||||
|
class MessageBase;
|
||||||
|
class MessageListenerBase;
|
||||||
|
|
||||||
|
class MessageQueue
|
||||||
|
{
|
||||||
|
public:
|
||||||
|
static std::shared_ptr<MessageQueue> getInstance();
|
||||||
|
|
||||||
|
void registerListener(MessageListenerBase* listener);
|
||||||
|
void unregisterListener(MessageListenerBase* listener);
|
||||||
|
|
||||||
|
void pushMessage(std::shared_ptr<MessageBase> message);
|
||||||
|
|
||||||
|
void startMessageLoopThreaded();
|
||||||
|
void startMessageLoop();
|
||||||
|
void stopMessageLoop();
|
||||||
|
|
||||||
|
bool loopIsRunning() const;
|
||||||
|
|
||||||
|
private:
|
||||||
|
typedef std::queue<std::shared_ptr<MessageBase> > MessageBufferType;
|
||||||
|
|
||||||
|
static std::shared_ptr<MessageQueue> s_instance;
|
||||||
|
|
||||||
|
MessageQueue();
|
||||||
|
MessageQueue(const MessageQueue&);
|
||||||
|
void operator=(const MessageQueue&);
|
||||||
|
|
||||||
|
void processMessages();
|
||||||
|
|
||||||
|
std::shared_ptr<MessageBufferType> m_frontMessageBuffer;
|
||||||
|
std::shared_ptr<MessageBufferType> m_backMessageBuffer;
|
||||||
|
std::vector<MessageListenerBase*> 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
|
||||||
@@ -4,6 +4,7 @@ add_files(
|
|||||||
ConfigManagerTestSuite.h
|
ConfigManagerTestSuite.h
|
||||||
CxxParserTestSuite.h
|
CxxParserTestSuite.h
|
||||||
LogManagerTestSuite.h
|
LogManagerTestSuite.h
|
||||||
|
MessageQueueTestSuite.h
|
||||||
TestSuiteFixture.cpp
|
TestSuiteFixture.cpp
|
||||||
TestSuiteFixture.h
|
TestSuiteFixture.h
|
||||||
Vector2TestSuite.h
|
Vector2TestSuite.h
|
||||||
|
|||||||
@@ -0,0 +1,203 @@
|
|||||||
|
#include <cxxtest/TestSuite.h>
|
||||||
|
|
||||||
|
#include <chrono>
|
||||||
|
|
||||||
|
#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<TestMessage>
|
||||||
|
{
|
||||||
|
public:
|
||||||
|
static const std::string getStaticType()
|
||||||
|
{
|
||||||
|
return "TestMessage";
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
class Test2Message: public Message<Test2Message>
|
||||||
|
{
|
||||||
|
public:
|
||||||
|
static const std::string getStaticType()
|
||||||
|
{
|
||||||
|
return "TestMessage2";
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
class TestMessageListener: public MessageListener<TestMessage>
|
||||||
|
{
|
||||||
|
public:
|
||||||
|
TestMessageListener()
|
||||||
|
: m_messageCount(0)
|
||||||
|
{
|
||||||
|
}
|
||||||
|
|
||||||
|
int m_messageCount;
|
||||||
|
|
||||||
|
private:
|
||||||
|
virtual void handleMessage(TestMessage* message)
|
||||||
|
{
|
||||||
|
m_messageCount++;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
class Test2MessageListener: public MessageListener<Test2Message>
|
||||||
|
{
|
||||||
|
public:
|
||||||
|
Test2MessageListener()
|
||||||
|
: m_messageCount(0)
|
||||||
|
{
|
||||||
|
}
|
||||||
|
|
||||||
|
int m_messageCount;
|
||||||
|
|
||||||
|
private:
|
||||||
|
virtual void handleMessage(Test2Message* message)
|
||||||
|
{
|
||||||
|
m_messageCount++;
|
||||||
|
TestMessage().dispatch();
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
class Test3MessageListener: public MessageListener<Test2Message>
|
||||||
|
{
|
||||||
|
public:
|
||||||
|
std::shared_ptr<TestMessageListener> m_listener;
|
||||||
|
|
||||||
|
private:
|
||||||
|
virtual void handleMessage(Test2Message* message)
|
||||||
|
{
|
||||||
|
m_listener = std::make_shared<TestMessageListener>();
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
class Test4MessageListener:
|
||||||
|
public MessageListener<TestMessage>,
|
||||||
|
public MessageListener<Test2Message>
|
||||||
|
{
|
||||||
|
public:
|
||||||
|
std::shared_ptr<TestMessageListener> m_listener;
|
||||||
|
|
||||||
|
private:
|
||||||
|
virtual void handleMessage(TestMessage* message)
|
||||||
|
{
|
||||||
|
if (!m_listener)
|
||||||
|
{
|
||||||
|
m_listener = std::make_shared<TestMessageListener>();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
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));
|
||||||
|
}
|
||||||
|
};
|
||||||
Reference in New Issue
Block a user