Fix for crash when unsubscription during event processing (untested)
This commit is contained in:
@@ -13,6 +13,8 @@
|
|||||||
relay = decltype(relay)(std::bind(handler, this, std::placeholders::_1)); \
|
relay = decltype(relay)(std::bind(handler, this, std::placeholders::_1)); \
|
||||||
m_EventBroker->Subscribe(relay);
|
m_EventBroker->Subscribe(relay);
|
||||||
|
|
||||||
|
typedef unsigned int EventID;
|
||||||
|
|
||||||
class EventBroker;
|
class EventBroker;
|
||||||
|
|
||||||
class BaseEventRelay
|
class BaseEventRelay
|
||||||
@@ -31,6 +33,7 @@ public:
|
|||||||
virtual bool Receive(const std::shared_ptr<Event> event) = 0;
|
virtual bool Receive(const std::shared_ptr<Event> event) = 0;
|
||||||
|
|
||||||
protected:
|
protected:
|
||||||
|
EventID m_EventID;
|
||||||
std::string m_ContextTypeName;
|
std::string m_ContextTypeName;
|
||||||
std::string m_EventTypeName;
|
std::string m_EventTypeName;
|
||||||
EventBroker* m_Broker;
|
EventBroker* m_Broker;
|
||||||
@@ -95,6 +98,7 @@ public:
|
|||||||
|
|
||||||
private:
|
private:
|
||||||
bool m_IsProcessing = false;
|
bool m_IsProcessing = false;
|
||||||
|
EventID m_NextEventID = 0;
|
||||||
|
|
||||||
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()
|
||||||
@@ -103,14 +107,14 @@ private:
|
|||||||
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_RelaysToSubscribe;
|
||||||
std::vector<BaseEventRelay*> m_RelaysToUnsubscribe;
|
std::vector<std::tuple<EventID, ContextTypeName_t, EventTypeName_t>> 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 subscribeImmediate(BaseEventRelay& relay);
|
||||||
void unsubscribeImmediate(BaseEventRelay& relay);
|
void unsubscribeImmediate(std::tuple<EventID, ContextTypeName_t, EventTypeName_t> identifier);
|
||||||
};
|
};
|
||||||
|
|
||||||
template <typename EventType>
|
template <typename EventType>
|
||||||
|
|||||||
@@ -2,21 +2,24 @@
|
|||||||
|
|
||||||
BaseEventRelay::~BaseEventRelay()
|
BaseEventRelay::~BaseEventRelay()
|
||||||
{
|
{
|
||||||
if (m_Broker != nullptr) {
|
if (m_Broker != nullptr) {
|
||||||
m_Broker->Unsubscribe(*this);
|
m_Broker->Unsubscribe(*this);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
void EventBroker::Unsubscribe(BaseEventRelay &relay) // ?
|
void EventBroker::Unsubscribe(BaseEventRelay& relay) // ?
|
||||||
{
|
{
|
||||||
if (m_IsProcessing) {
|
auto identifier = std::make_tuple(relay.m_EventID, relay.m_ContextTypeName, relay.m_EventTypeName);
|
||||||
m_RelaysToUnsubscribe.push_back(&relay);
|
|
||||||
} else {
|
relay.m_Broker = nullptr;
|
||||||
unsubscribeImmediate(relay);
|
if (m_IsProcessing) {
|
||||||
}
|
m_RelaysToUnsubscribe.push_back(identifier);
|
||||||
|
} else {
|
||||||
|
unsubscribeImmediate(identifier);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
void EventBroker::Subscribe(BaseEventRelay &relay)
|
void EventBroker::Subscribe(BaseEventRelay& relay)
|
||||||
{
|
{
|
||||||
if (m_IsProcessing) {
|
if (m_IsProcessing) {
|
||||||
m_RelaysToSubscribe.push_back(&relay);
|
m_RelaysToSubscribe.push_back(&relay);
|
||||||
@@ -38,12 +41,11 @@ int EventBroker::Process(std::string contextTypeName)
|
|||||||
|
|
||||||
int eventsProcessed = 0;
|
int eventsProcessed = 0;
|
||||||
for (auto &pair : *m_EventQueueRead) {
|
for (auto &pair : *m_EventQueueRead) {
|
||||||
std::string &eventTypeName = pair.first;
|
std::string& eventTypeName = pair.first;
|
||||||
std::shared_ptr<Event> event = pair.second;
|
std::shared_ptr<Event> event = pair.second;
|
||||||
|
|
||||||
auto itpair = relays.equal_range(eventTypeName);
|
auto itpair = relays.equal_range(eventTypeName);
|
||||||
for (auto it2 = itpair.first; it2 != itpair.second; it2++)
|
for (auto it2 = itpair.first; it2 != itpair.second; it2++) {
|
||||||
{
|
|
||||||
std::string name = it2->first;
|
std::string name = it2->first;
|
||||||
BaseEventRelay* relay = it2->second;
|
BaseEventRelay* relay = it2->second;
|
||||||
relay->Receive(event);
|
relay->Receive(event);
|
||||||
@@ -60,8 +62,8 @@ int EventBroker::Process(std::string contextTypeName)
|
|||||||
m_RelaysToSubscribe.clear();
|
m_RelaysToSubscribe.clear();
|
||||||
|
|
||||||
// Process pending unsubscriptions
|
// Process pending unsubscriptions
|
||||||
for (auto& r : m_RelaysToUnsubscribe) {
|
for (auto& identifier : m_RelaysToUnsubscribe) {
|
||||||
unsubscribeImmediate(*r);
|
unsubscribeImmediate(identifier);
|
||||||
}
|
}
|
||||||
m_RelaysToUnsubscribe.clear();
|
m_RelaysToUnsubscribe.clear();
|
||||||
|
|
||||||
@@ -81,21 +83,26 @@ void EventBroker::Clear()
|
|||||||
void EventBroker::subscribeImmediate(BaseEventRelay& relay)
|
void EventBroker::subscribeImmediate(BaseEventRelay& relay)
|
||||||
{
|
{
|
||||||
relay.m_Broker = this;
|
relay.m_Broker = this;
|
||||||
|
relay.m_EventID = m_NextEventID++;
|
||||||
m_ContextRelays[relay.m_ContextTypeName].insert(std::make_pair(relay.m_EventTypeName, &relay));
|
m_ContextRelays[relay.m_ContextTypeName].insert(std::make_pair(relay.m_EventTypeName, &relay));
|
||||||
}
|
}
|
||||||
|
|
||||||
void EventBroker::unsubscribeImmediate(BaseEventRelay& relay)
|
void EventBroker::unsubscribeImmediate(std::tuple<EventID, ContextTypeName_t, EventTypeName_t> identifier)
|
||||||
{
|
{
|
||||||
auto contextIt = m_ContextRelays.find(relay.m_ContextTypeName);
|
EventID eventID;
|
||||||
|
ContextTypeName_t contextTypeName;
|
||||||
|
EventTypeName_t eventTypeName;
|
||||||
|
std::tie(eventID, contextTypeName, eventTypeName) = identifier;
|
||||||
|
|
||||||
|
auto contextIt = m_ContextRelays.find(contextTypeName);
|
||||||
if (contextIt == m_ContextRelays.end()) {
|
if (contextIt == m_ContextRelays.end()) {
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
auto eventRelays = contextIt->second;
|
auto eventRelays = contextIt->second;
|
||||||
auto itpair = eventRelays.equal_range(relay.m_EventTypeName);
|
auto itpair = eventRelays.equal_range(eventTypeName);
|
||||||
for (auto it = itpair.first; it != itpair.second; ++it) {
|
for (auto it = itpair.first; it != itpair.second; ++it) {
|
||||||
if (it->second == &relay) {
|
if (it->second->m_EventID == eventID) {
|
||||||
relay.m_Broker = nullptr;
|
|
||||||
eventRelays.erase(it);
|
eventRelays.erase(it);
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user