Merge branch 'master' of github.com:john30/ebusd

This commit is contained in:
john30
2018-01-28 14:12:05 +01:00
2 changed files with 32 additions and 11 deletions
+29 -9
View File
@@ -301,7 +301,7 @@ void on_connect(
MqttHandler::MqttHandler(UserInfo* userInfo, BusHandler* busHandler, MessageMap* messages) MqttHandler::MqttHandler(UserInfo* userInfo, BusHandler* busHandler, MessageMap* messages)
: DataSink(userInfo, "mqtt"), DataSource(busHandler), Thread(), m_messages(messages), m_connected(false), : DataSink(userInfo, "mqtt"), DataSource(busHandler), WaitThread(), m_messages(messages), m_connected(false),
m_lastUpdateCheckResult(".") { m_lastUpdateCheckResult(".") {
m_publishByField = false; m_publishByField = false;
m_mosquitto = NULL; m_mosquitto = NULL;
@@ -375,14 +375,25 @@ MqttHandler::MqttHandler(UserInfo* userInfo, BusHandler* busHandler, MessageMap*
#endif #endif
mosquitto_connect_callback_set(m_mosquitto, on_connect); mosquitto_connect_callback_set(m_mosquitto, on_connect);
int ret;
#if (LIBMOSQUITTO_MAJOR >= 1) #if (LIBMOSQUITTO_MAJOR >= 1)
if (mosquitto_connect(m_mosquitto, g_host, g_port, 60) != MOSQ_ERR_SUCCESS) { ret = mosquitto_connect(m_mosquitto, g_host, g_port, 60);
#else #else
if (mosquitto_connect(m_mosquitto, g_host, g_port, 60, true) != MOSQ_ERR_SUCCESS) { ret = mosquitto_connect(m_mosquitto, g_host, g_port, 60, true);
#endif #endif
logOtherError("mqtt", "unable to connect"); if (ret == MOSQ_ERR_INVAL) {
logOtherError("mqtt", "unable to connect (invalid parameters)");
mosquitto_destroy(m_mosquitto); mosquitto_destroy(m_mosquitto);
m_mosquitto = NULL; m_mosquitto = NULL;
} else if (ret != MOSQ_ERR_SUCCESS) {
m_connected = false;
char* error;
if (ret != MOSQ_ERR_ERRNO) {
error = (char*)("unkown mosquitto error");
} else {
error = strerror(errno);
}
logOtherError("mqtt", "unable to connect, retrying: %s", error);
} else { } else {
m_connected = true; // assume success until connect_callback says otherwise m_connected = true; // assume success until connect_callback says otherwise
logOtherDebug("mqtt", "connection requested"); logOtherDebug("mqtt", "connection requested");
@@ -401,7 +412,7 @@ MqttHandler::~MqttHandler() {
void MqttHandler::start() { void MqttHandler::start() {
if (m_mosquitto) { if (m_mosquitto) {
Thread::start("MQTT"); WaitThread::start("MQTT");
} }
} }
@@ -526,8 +537,10 @@ void MqttHandler::run() {
mosquitto_message_callback_set(m_mosquitto, on_message); mosquitto_message_callback_set(m_mosquitto, on_message);
string subTopic = getTopic(NULL, "#"); string subTopic = getTopic(NULL, "#");
mosquitto_subscribe(m_mosquitto, NULL, subTopic.c_str(), 0); mosquitto_subscribe(m_mosquitto, NULL, subTopic.c_str(), 0);
bool allowReconnect = false;
while (isRunning()) { while (isRunning()) {
handleTraffic(); handleTraffic(allowReconnect);
allowReconnect = false;
time(&now); time(&now);
if (now < start) { if (now < start) {
// clock skew // clock skew
@@ -536,6 +549,7 @@ void MqttHandler::run() {
} }
lastTaskRun = now; lastTaskRun = now;
} else if (now > lastTaskRun+15) { } else if (now > lastTaskRun+15) {
allowReconnect = true;
if (m_busHandler->hasSignal()) { if (m_busHandler->hasSignal()) {
lastSignal = now; lastSignal = now;
if (!signal) { if (!signal) {
@@ -577,17 +591,23 @@ void MqttHandler::run() {
m_updatedMessages.clear(); m_updatedMessages.clear();
} }
} }
if (!m_connected && !Wait(5)) {
break;
}
} }
} }
void MqttHandler::handleTraffic() { void MqttHandler::handleTraffic(bool allowReconnect) {
if (m_mosquitto) { if (m_mosquitto) {
int ret; int ret;
#if (LIBMOSQUITTO_MAJOR >= 1) #if (LIBMOSQUITTO_MAJOR >= 1)
ret = mosquitto_loop(m_mosquitto, -1, 1); ret = mosquitto_loop(m_mosquitto, -1, 1); // waits up to 1 second for network traffic
#else #else
ret = mosquitto_loop(m_mosquitto, -1); ret = mosquitto_loop(m_mosquitto, -1); // waits up to 1 second for network traffic
#endif #endif
if (!m_connected && ret == MOSQ_ERR_NO_CONN && allowReconnect) {
ret = mosquitto_reconnect(m_mosquitto);
}
if (!m_connected && ret == MOSQ_ERR_SUCCESS) { if (!m_connected && ret == MOSQ_ERR_SUCCESS) {
m_connected = true; m_connected = true;
logOtherNotice("mqtt", "connection re-established"); logOtherNotice("mqtt", "connection re-established");
+3 -2
View File
@@ -58,7 +58,7 @@ bool mqtthandler_register(UserInfo* userInfo, BusHandler* busHandler, MessageMap
/** /**
* The main class supporting MQTT data handling. * The main class supporting MQTT data handling.
*/ */
class MqttHandler : public DataSink, public DataSource, public Thread { class MqttHandler : public DataSink, public DataSource, public WaitThread {
public: public:
/** /**
* Constructor. * Constructor.
@@ -94,8 +94,9 @@ class MqttHandler : public DataSink, public DataSource, public Thread {
private: private:
/** /**
* Called regularly to handle MQTT traffic. * Called regularly to handle MQTT traffic.
* @param allowReconnect true when reconnecting to the broker is allowed.
*/ */
void handleTraffic(); void handleTraffic(bool allowReconnect);
/** /**
* Build the MQTT topic string for the @a Message. * Build the MQTT topic string for the @a Message.