Fixed crash when events were (un)subscribed while events were processing

This commit is contained in:
2015-10-02 18:12:00 +02:00
parent 45f1d882ad
commit e74f1092c7
3 changed files with 59 additions and 18 deletions
+8
View File
@@ -24,6 +24,7 @@
#include <map> #include <map>
#include <unordered_map> #include <unordered_map>
#include <list> #include <list>
#include <tuple>
#define EVENT_SUBSCRIBE_MEMBER(relay, handler) \ #define EVENT_SUBSCRIBE_MEMBER(relay, handler) \
relay = decltype(relay)(std::bind(handler, this, std::placeholders::_1)); \ relay = decltype(relay)(std::bind(handler, this, std::placeholders::_1)); \
@@ -123,16 +124,23 @@ public:
void Unsubscribe(BaseEventRelay &relay); void Unsubscribe(BaseEventRelay &relay);
private: private:
bool m_IsProcessing = false;
typedef std::string ContextTypeName_t; // typeid(ContextType).name() typedef std::string ContextTypeName_t; // typeid(ContextType).name()
typedef std::string EventTypeName_t; // typeid(EventType).name() typedef std::string EventTypeName_t; // typeid(EventType).name()
typedef std::unordered_multimap<EventTypeName_t, BaseEventRelay*> EventRelays_t; typedef std::unordered_multimap<EventTypeName_t, BaseEventRelay*> EventRelays_t;
typedef std::unordered_map<ContextTypeName_t, EventRelays_t> ContextRelays_t; typedef std::unordered_map<ContextTypeName_t, EventRelays_t> ContextRelays_t;
ContextRelays_t m_ContextRelays; ContextRelays_t m_ContextRelays;
std::vector<BaseEventRelay*> m_RelaysToSubscribe;
std::vector<BaseEventRelay*> m_RelaysToUnsubscribe;
typedef std::list<std::pair<EventTypeName_t, std::shared_ptr<Event>>> EventQueue_t; typedef std::list<std::pair<EventTypeName_t, std::shared_ptr<Event>>> EventQueue_t;
std::shared_ptr<EventQueue_t> m_EventQueueRead; std::shared_ptr<EventQueue_t> m_EventQueueRead;
std::shared_ptr<EventQueue_t> m_EventQueueWrite; std::shared_ptr<EventQueue_t> m_EventQueueWrite;
void subscribeImmediate(BaseEventRelay& relay);
void unsubscribeImmediate(BaseEventRelay& relay);
}; };
template <typename EventType> template <typename EventType>
+50 -17
View File
@@ -29,32 +29,26 @@ dd::BaseEventRelay::~BaseEventRelay()
void dd::EventBroker::Unsubscribe(BaseEventRelay &relay) // ? void dd::EventBroker::Unsubscribe(BaseEventRelay &relay) // ?
{ {
auto contextIt = m_ContextRelays.find(relay.m_ContextTypeName); if (m_IsProcessing) {
if (contextIt == m_ContextRelays.end()) m_RelaysToUnsubscribe.push_back(&relay);
return; } else {
unsubscribeImmediate(relay);
auto eventRelays = contextIt->second;
auto itpair = eventRelays.equal_range(relay.m_EventTypeName);
for (auto it = itpair.first; it != itpair.second; ++it)
{
if (it->second == &relay)
{
relay.m_Broker = nullptr;
eventRelays.erase(it);
break;
}
} }
} }
void dd::EventBroker::Subscribe(BaseEventRelay &relay) void dd::EventBroker::Subscribe(BaseEventRelay &relay)
{ {
relay.m_Broker = this; if (m_IsProcessing) {
m_ContextRelays[relay.m_ContextTypeName].insert(std::make_pair(relay.m_EventTypeName, &relay)); m_RelaysToSubscribe.push_back(&relay);
} else {
subscribeImmediate(relay);
}
} }
int dd::EventBroker::Process(std::string contextTypeName) int dd::EventBroker::Process(std::string contextTypeName)
{ {
m_IsProcessing = true;
auto it = m_ContextRelays.find(contextTypeName); auto it = m_ContextRelays.find(contextTypeName);
if (it == m_ContextRelays.end()) if (it == m_ContextRelays.end())
return 0; return 0;
@@ -77,6 +71,20 @@ int dd::EventBroker::Process(std::string contextTypeName)
} }
} }
m_IsProcessing = false;
// Process pending subscriptions
for (auto& r : m_RelaysToSubscribe) {
subscribeImmediate(*r);
}
m_RelaysToSubscribe.clear();
// Process pending unsubscriptions
for (auto& r : m_RelaysToUnsubscribe) {
unsubscribeImmediate(*r);
}
m_RelaysToUnsubscribe.clear();
return eventsProcessed; return eventsProcessed;
} }
@@ -84,4 +92,29 @@ void dd::EventBroker::Swap()
{ {
std::swap(m_EventQueueRead, m_EventQueueWrite); std::swap(m_EventQueueRead, m_EventQueueWrite);
m_EventQueueWrite->clear(); m_EventQueueWrite->clear();
}
void dd::EventBroker::subscribeImmediate(dd::BaseEventRelay& relay)
{
relay.m_Broker = this;
m_ContextRelays[relay.m_ContextTypeName].insert(std::make_pair(relay.m_EventTypeName, &relay));
}
void dd::EventBroker::unsubscribeImmediate(dd::BaseEventRelay& relay)
{
auto contextIt = m_ContextRelays.find(relay.m_ContextTypeName);
if (contextIt == m_ContextRelays.end()) {
return;
}
auto eventRelays = contextIt->second;
auto itpair = eventRelays.equal_range(relay.m_EventTypeName);
for (auto it = itpair.first; it != itpair.second; ++it) {
if (it->second == &relay) {
relay.m_Broker = nullptr;
eventRelays.erase(it);
break;
}
}
} }
+1 -1
View File
@@ -10,7 +10,7 @@ struct EventFixture
EventFixture() EventFixture()
{ {
this->EventBroker = new dd::EventBroker(); this->EventBroker = new dd::EventBroker();
m_EEventType = decltype(m_EEventType)(std::bind(&OnEvent, this, std::placeholders::_1)); m_EEventType = decltype(m_EEventType)(std::bind(&EventFixture::OnEvent, this, std::placeholders::_1));
this->EventBroker->Subscribe(m_EEventType); this->EventBroker->Subscribe(m_EEventType);
Run(); Run();
Check(); Check();