MqttClient was unable to publish in some cases

This commit is contained in:
hsaturn
2021-03-22 01:19:26 +01:00
parent 620dbf31af
commit e71a4d5e87
2 changed files with 28 additions and 16 deletions

View File

@@ -141,8 +141,10 @@ void MqttBroker::loop()
} }
} }
void MqttBroker::publish(const MqttClient* source, const Topic& topic, MqttMessage& msg) MqttError MqttBroker::publish(const MqttClient* source, const Topic& topic, MqttMessage& msg)
{ {
MqttError retval = MqttOk;
debug("publish "); debug("publish ");
int i=0; int i=0;
for(auto client: clients) for(auto client: clients)
@@ -158,7 +160,7 @@ void MqttBroker::publish(const MqttClient* source, const Topic& topic, MqttMessa
if (source == broker) // broker -> clients if (source == broker) // broker -> clients
doit = true; doit = true;
else // clients -> broker else // clients -> broker
broker->publish(topic, msg); retval=broker->publish(topic, msg);
} }
else // Disconnected: R7 else // Disconnected: R7
{ {
@@ -170,6 +172,7 @@ void MqttBroker::publish(const MqttClient* source, const Topic& topic, MqttMessa
if (doit) client->publish(topic, msg); if (doit) client->publish(topic, msg);
debug(""); debug("");
} }
return retval;
} }
bool MqttBroker::compareString( bool MqttBroker::compareString(
@@ -390,22 +393,25 @@ bool Topic::matches(const Topic& topic) const
} }
// publish from local client // publish from local client
void MqttClient::publish(const Topic& topic, const char* payload, size_t pay_length) MqttError MqttClient::publish(const Topic& topic, const char* payload, size_t pay_length)
{ {
message.create(MqttMessage::Publish); MqttMessage msg;
message.add(topic); msg.create(MqttMessage::Publish);
message.add(payload, pay_length); msg.add(topic);
msg.add(payload, pay_length);
if (parent) if (parent)
parent->publish(this, topic, message); return parent->publish(this, topic, msg);
else if (client) else if (client)
publish(topic, message); msg.sendTo(this);
else else
Serial << " Should not happen" << endl; return MqttNowhereToSend;
} }
// republish a received publish if it matches any in subscriptions // republish a received publish if it matches any in subscriptions
void MqttClient::publish(const Topic& topic, MqttMessage& msg) MqttError MqttClient::publish(const Topic& topic, MqttMessage& msg)
{ {
MqttError retval=MqttOk;
debug("mqttclient publish " << subscriptions.size()); debug("mqttclient publish " << subscriptions.size());
for(const auto& subscription: subscriptions) for(const auto& subscription: subscriptions)
{ {
@@ -419,11 +425,12 @@ void MqttClient::publish(const Topic& topic, MqttMessage& msg)
} }
else if (callback) else if (callback)
{ {
callback(this, topic, nullptr, 0); // TODO callback(this, topic, nullptr, 0); // TODO Payload
} }
} }
Serial << endl; Serial << endl;
} }
return retval;
} }
void MqttMessage::reset() void MqttMessage::reset()

View File

@@ -4,6 +4,11 @@
#include <string> #include <string>
#include "StringIndexer.h" #include "StringIndexer.h"
enum MqttError
{
MqttOk = 0,
MqttNowhereToSend=1,
};
class Topic : public IndexedString class Topic : public IndexedString
{ {
@@ -121,9 +126,9 @@ class MqttClient
void setCallback(CallBack fun) {callback=fun; }; void setCallback(CallBack fun) {callback=fun; };
// Publish from client to the world // Publish from client to the world
void publish(const Topic&, const char* payload, size_t pay_length); MqttError publish(const Topic&, const char* payload, size_t pay_length);
void publish(const Topic& t, const std::string& s) { publish(t,s.c_str(),s.length());} MqttError publish(const Topic& t, const std::string& s) { return publish(t,s.c_str(),s.length());}
void publish(const Topic& t) { publish(t, nullptr, 0);}; MqttError publish(const Topic& t) { return publish(t, nullptr, 0);};
void subscribe(Topic topic) { subscriptions.insert(topic); } void subscribe(Topic topic) { subscriptions.insert(topic); }
void unsubscribe(Topic& topic); void unsubscribe(Topic& topic);
@@ -151,7 +156,7 @@ class MqttClient
friend class MqttBroker; friend class MqttBroker;
MqttClient(MqttBroker* parent, WiFiClient& client); MqttClient(MqttBroker* parent, WiFiClient& client);
// republish a received publish if topic matches any in subscriptions // republish a received publish if topic matches any in subscriptions
void publish(const Topic& topic, MqttMessage& msg); MqttError publish(const Topic& topic, MqttMessage& msg);
void clientAlive(uint32_t more_seconds); void clientAlive(uint32_t more_seconds);
void processMessage(); void processMessage();
@@ -214,7 +219,7 @@ class MqttBroker
{ return compareString(auth_password, password, len); } { return compareString(auth_password, password, len); }
void publish(const MqttClient* source, const Topic& topic, MqttMessage& msg); MqttError publish(const MqttClient* source, const Topic& topic, MqttMessage& msg);
// For clients that are added not by the broker itself // For clients that are added not by the broker itself
void addClient(MqttClient* client); void addClient(MqttClient* client);