Return [numSubscribers, success] when (un)subscribing
This commit is contained in:
+10
-6
@@ -278,7 +278,8 @@ public:
|
|||||||
return emptyVector;
|
return emptyVector;
|
||||||
}
|
}
|
||||||
|
|
||||||
void subscribe(std::string_view topic, Subscriber *subscriber, bool nonStrict = false) {
|
/* Returns number of subscribers after the call and whether or not we were successful in subscribing */
|
||||||
|
std::pair<unsigned int, bool> subscribe(std::string_view topic, Subscriber *subscriber, bool nonStrict = false) {
|
||||||
/* Start iterating from the root */
|
/* Start iterating from the root */
|
||||||
Topic *iterator = root;
|
Topic *iterator = root;
|
||||||
|
|
||||||
@@ -333,7 +334,9 @@ public:
|
|||||||
/* Add Topic to list of subscriptions only if we weren't already subscribed */
|
/* Add Topic to list of subscriptions only if we weren't already subscribed */
|
||||||
if (inserted) {
|
if (inserted) {
|
||||||
subscriber->subscriptions.push_back(iterator);
|
subscriber->subscriptions.push_back(iterator);
|
||||||
|
return {(unsigned int) iterator->subs.size(), true};
|
||||||
}
|
}
|
||||||
|
return {(unsigned int) iterator->subs.size(), false};
|
||||||
}
|
}
|
||||||
|
|
||||||
void publish(std::string_view topic, std::pair<std::string_view, std::string_view> message, Subscriber *sender = nullptr) {
|
void publish(std::string_view topic, std::pair<std::string_view, std::string_view> message, Subscriber *sender = nullptr) {
|
||||||
@@ -347,8 +350,8 @@ public:
|
|||||||
messageId++;
|
messageId++;
|
||||||
}
|
}
|
||||||
|
|
||||||
/* Returns whether we were subscribed prior */
|
/* Returns a pair of numSubscribers after operation, and whether we were subscribed prior */
|
||||||
bool unsubscribe(std::string_view topic, Subscriber *subscriber, bool nonStrict = false) {
|
std::pair<unsigned int, bool> unsubscribe(std::string_view topic, Subscriber *subscriber, bool nonStrict = false) {
|
||||||
/* Subscribers are likely to have very few subscriptions (20 or fewer) */
|
/* Subscribers are likely to have very few subscriptions (20 or fewer) */
|
||||||
if (subscriber) {
|
if (subscriber) {
|
||||||
/* Lookup exact Topic ptr from string */
|
/* Lookup exact Topic ptr from string */
|
||||||
@@ -360,7 +363,7 @@ public:
|
|||||||
std::map<std::string_view, Topic *>::iterator it = iterator->children.find(segment);
|
std::map<std::string_view, Topic *>::iterator it = iterator->children.find(segment);
|
||||||
if (it == iterator->children.end()) {
|
if (it == iterator->children.end()) {
|
||||||
/* This topic does not even exist */
|
/* This topic does not even exist */
|
||||||
return false;
|
return {0, false};
|
||||||
}
|
}
|
||||||
|
|
||||||
iterator = it->second;
|
iterator = it->second;
|
||||||
@@ -381,12 +384,13 @@ public:
|
|||||||
|
|
||||||
/* Remove us from Topic's subs */
|
/* Remove us from Topic's subs */
|
||||||
iterator->subs.erase(subscriber);
|
iterator->subs.erase(subscriber);
|
||||||
|
unsigned int numSubscribers = (unsigned int) iterator->subs.size();
|
||||||
trimTree(iterator);
|
trimTree(iterator);
|
||||||
return true;
|
return {numSubscribers, true};
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return false;
|
return {0, false};
|
||||||
}
|
}
|
||||||
|
|
||||||
/* Can be called with nullptr, ignore it then */
|
/* Can be called with nullptr, ignore it then */
|
||||||
|
|||||||
+5
-5
@@ -197,8 +197,8 @@ public:
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/* Subscribe to a topic according to MQTT rules and syntax */
|
/* Subscribe to a topic according to MQTT rules and syntax. Returns [numSubscribers, success]. */
|
||||||
void subscribe(std::string_view topic, bool nonStrict = false) {
|
std::pair<unsigned int, bool> subscribe(std::string_view topic, bool nonStrict = false) {
|
||||||
WebSocketContextData<SSL> *webSocketContextData = (WebSocketContextData<SSL> *) us_socket_context_ext(SSL,
|
WebSocketContextData<SSL> *webSocketContextData = (WebSocketContextData<SSL> *) us_socket_context_ext(SSL,
|
||||||
(us_socket_context_t *) us_socket_context(SSL, (us_socket_t *) this)
|
(us_socket_context_t *) us_socket_context(SSL, (us_socket_t *) this)
|
||||||
);
|
);
|
||||||
@@ -209,11 +209,11 @@ public:
|
|||||||
webSocketData->subscriber = new Subscriber(this);
|
webSocketData->subscriber = new Subscriber(this);
|
||||||
}
|
}
|
||||||
|
|
||||||
webSocketContextData->topicTree.subscribe(topic, webSocketData->subscriber, nonStrict);
|
return webSocketContextData->topicTree.subscribe(topic, webSocketData->subscriber, nonStrict);
|
||||||
}
|
}
|
||||||
|
|
||||||
/* Unsubscribe from a topic, returns true if we were subscribed */
|
/* Unsubscribe from a topic, returns true if we were subscribed. Returns [numSubscribers, success]. */
|
||||||
bool unsubscribe(std::string_view topic, bool nonStrict = false) {
|
std::pair<unsigned int, bool> unsubscribe(std::string_view topic, bool nonStrict = false) {
|
||||||
WebSocketContextData<SSL> *webSocketContextData = (WebSocketContextData<SSL> *) us_socket_context_ext(SSL,
|
WebSocketContextData<SSL> *webSocketContextData = (WebSocketContextData<SSL> *) us_socket_context_ext(SSL,
|
||||||
(us_socket_context_t *) us_socket_context(SSL, (us_socket_t *) this)
|
(us_socket_context_t *) us_socket_context(SSL, (us_socket_t *) this)
|
||||||
);
|
);
|
||||||
|
|||||||
Reference in New Issue
Block a user