diff --git a/src/TopicTreeDraft.h b/src/TopicTreeDraft.h index 85e53c8..0f832ad 100644 --- a/src/TopicTreeDraft.h +++ b/src/TopicTreeDraft.h @@ -120,6 +120,11 @@ private: /* Should be getData and commit? */ void publish(Topic *iterator, size_t start, size_t stop, std::string_view topic, std::string_view message) { + /* If we already have 64 triggered topics make sure to drain it here */ + if (numTriggeredTopics == 64) { + drain(); + } + for (; stop != std::string::npos; start = stop + 1) { stop = topic.find('/', start); std::string_view segment = topic.substr(start, stop - start); @@ -235,10 +240,39 @@ public: messageId++; } - /* Rarely used, probably */ - void unsubscribe(std::string_view topic, Subscriber *subscriber) { + /* Returns whether we were subscribed prior */ + bool unsubscribe(std::string_view topic, Subscriber *subscriber) { /* Subscribers are likely to have very few subscriptions (20 or fewer) */ + if (subscriber) { + /* Lookup exact Topic ptr from string */ + Topic *iterator = root; + for (size_t start = 0, stop = 0; stop != std::string::npos; start = stop + 1) { + stop = topic.find('/', start); + std::string_view segment = topic.substr(start, stop - start); + std::map::iterator it = iterator->children.find(segment); + if (it == iterator->children.end()) { + /* This topic does not even exist */ + return false; + } + + iterator = it->second; + } + + /* Try and remove this topic from our list */ + for (auto it = subscriber->subscriptions.begin(); it != subscriber->subscriptions.end(); it++) { + if (*it == iterator) { + /* Remove topic ptr from our list */ + subscriber->subscriptions.erase(it); + + /* Remove us from Topic's subs */ + iterator->subs.erase(subscriber); + trimTree(iterator); + return true; + } + } + } + return false; } /* Can be called with nullptr, ignore it then */ diff --git a/src/WebSocket.h b/src/WebSocket.h index 83c5c60..9fb322c 100644 --- a/src/WebSocket.h +++ b/src/WebSocket.h @@ -154,6 +154,17 @@ public: webSocketContextData->topicTree.subscribe(topic, webSocketData->subscriber); } + /* Unsubscribe from a topic, returns true if we were subscribed */ + bool unsubscribe(std::string_view topic) { + WebSocketContextData *webSocketContextData = (WebSocketContextData *) us_socket_context_ext(SSL, + (us_socket_context_t *) us_socket_context(SSL, (us_socket_t *) this) + ); + + WebSocketData *webSocketData = (WebSocketData *) us_socket_ext(SSL, (us_socket_t *) this); + + return webSocketContextData->topicTree.unsubscribe(topic, webSocketData->subscriber); + } + /* Publish a message to a topic according to MQTT rules and syntax */ void publish(std::string_view topic, std::string_view message, OpCode opCode = OpCode::TEXT, bool compress = false) { WebSocketContextData *webSocketContextData = (WebSocketContextData *) us_socket_context_ext(SSL,