From cb78d981857d4315bfb3568ffcab1c0b1452f534 Mon Sep 17 00:00:00 2001 From: john30 Date: Tue, 25 Dec 2018 13:40:16 +0100 Subject: [PATCH] added --mqttlog and --mqttversion options, check return code of calls to mosquitto library, log library version, set threaded mode --- src/ebusd/mqtthandler.cpp | 178 +++++++++++++++++++++++++++----------- 1 file changed, 129 insertions(+), 49 deletions(-) diff --git a/src/ebusd/mqtthandler.cpp b/src/ebusd/mqtthandler.cpp index b2223232..d701f0eb 100755 --- a/src/ebusd/mqtthandler.cpp +++ b/src/ebusd/mqtthandler.cpp @@ -36,7 +36,9 @@ using std::dec; #define O_TOPI (O_PASS+1) #define O_RETA (O_TOPI+1) #define O_JSON (O_RETA+1) -#define O_IGIN (O_JSON+1) +#define O_LOGL (O_JSON+1) +#define O_VERS (O_LOGL+1) +#define O_IGIN (O_VERS+1) #define O_CHGS (O_IGIN+1) #define O_CAFI (O_CHGS+1) #define O_CERT (O_CAFI+1) @@ -56,6 +58,12 @@ static const struct argp_option g_mqtt_argp_options[] = { "Use MQTT TOPIC (prefix before /%circuit/%name or complete format) [ebusd]", 0 }, {"mqttretain", O_RETA, nullptr, 0, "Retain all topics instead of only selected global ones", 0 }, {"mqttjson", O_JSON, nullptr, 0, "Publish in JSON format instead of strings", 0 }, +#if (LIBMOSQUITTO_VERSION_NUMBER >= 1003001) + {"mqttlog", O_LOGL, nullptr, 0, "Log library events", 0 }, +#endif +#if (LIBMOSQUITTO_VERSION_NUMBER >= 1004001) + {"mqttversion", O_VERS, "VERSION", 0, "Use protocol VERSION [3.1]", 0 }, +#endif {"mqttignoreinvalid", O_IGIN, nullptr, 0, "Ignore invalid parameters during init (e.g. for DNS not resolvable yet)", 0 }, {"mqttchanges", O_CHGS, nullptr, 0, "Whether to only publish changed messages instead of all received", 0 }, @@ -81,6 +89,12 @@ static vector g_topicStrs; static vector g_topicFields; static bool g_retain = false; //!< whether to retail all topics static OutputFormat g_publishFormat = 0; //!< the OutputFormat for publishing messages +#if (LIBMOSQUITTO_VERSION_NUMBER >= 1003001) +static bool g_logFromLib = false; //!< log library events +#endif +#if (LIBMOSQUITTO_VERSION_NUMBER >= 1004001) +static int g_version = MQTT_PROTOCOL_V31; //!< protocol version to use +#endif static bool g_ignoreInvalidParams = false; //!< ignore invalid parameters during init static bool g_onlyChanges = true; //!< whether to only publish changed messages instead of all received @@ -162,6 +176,22 @@ static error_t mqtt_parse_opt(int key, char *arg, struct argp_state *state) { g_publishFormat |= OF_JSON|OF_NAMES; break; +#if (LIBMOSQUITTO_VERSION_NUMBER >= 1003001) + case O_LOGL: + g_logFromLib = true; + break; +#endif + +#if (LIBMOSQUITTO_VERSION_NUMBER >= 1004001) + case O_VERS: // --mqttversion=3.1.1 + if (arg == nullptr || arg[0] == 0 || (strcmp(arg, "3.1")!=0 && strcmp(arg, "3.1.1")!=0)) { + argp_error(state, "invalid mqttversion"); + return EINVAL; + } + g_version = strcmp(arg, "3.1.1")==0 ? MQTT_PROTOCOL_V311 : MQTT_PROTOCOL_V31; + break; +#endif + case O_IGIN: g_ignoreInvalidParams = true; break; @@ -228,15 +258,37 @@ const struct argp_child* mqtthandler_getargs() { return &g_mqtt_argp_child; } +bool check(int code, const char* method) { + if (code == MOSQ_ERR_SUCCESS) { + return true; + } + if (code == MOSQ_ERR_ERRNO) { + char* error = strerror(errno); + logOtherError("mqtt", "%s: errno %d=%s", method, errno, error); + return false; + } +#if (LIBMOSQUITTO_VERSION_NUMBER >= 1003001) + const char* msg = mosquitto_strerror(code); + logOtherError("mqtt", "%s: %s", method, msg); +#else + logOtherError("mqtt", "%s: error code %d", method, code); +#endif + return false; +} + bool mqtthandler_register(UserInfo* userInfo, BusHandler* busHandler, MessageMap* messages, list* handlers) { if (g_port > 0) { int major = -1; - mosquitto_lib_version(&major, nullptr, nullptr); + int minor = -1; + int revision = -1; + mosquitto_lib_version(&major, &minor, &revision); if (major != LIBMOSQUITTO_MAJOR) { logOtherError("mqtt", "invalid mosquitto version %d instead of %d", major, LIBMOSQUITTO_MAJOR); return false; } + logOtherInfo("mqtt", "mosquitto version %d.%d.%d (compiled with %d.%d.%d)", major, minor, revision, + LIBMOSQUITTO_MAJOR, LIBMOSQUITTO_MINOR, LIBMOSQUITTO_REVISION); handlers->push_back(new MqttHandler(userInfo, busHandler, messages)); } return true; @@ -331,6 +383,31 @@ void on_connect( } } +#if (LIBMOSQUITTO_VERSION_NUMBER >= 1003001) +void on_log(struct mosquitto *mosq, void *obj, int level, const char* msg) { + switch (level) { + case MOSQ_LOG_DEBUG: + logOtherDebug("mqtt", "log %s", msg); + break; + case MOSQ_LOG_INFO: + logOtherInfo("mqtt", "log %s", msg); + break; + case MOSQ_LOG_NOTICE: + logOtherNotice("mqtt", "log %s", msg); + break; + case MOSQ_LOG_WARNING: + logOtherNotice("mqtt", "log warning %s", msg); + break; + case MOSQ_LOG_ERR: + logOtherError("mqtt", "log %s", msg); + break; + default: + logOtherError("mqtt", "log other %s", msg); + break; + } +} +#endif + void on_message( #if (LIBMOSQUITTO_MAJOR >= 1) struct mosquitto *mosq, @@ -373,10 +450,7 @@ MqttHandler::MqttHandler(UserInfo* userInfo, BusHandler* busHandler, MessageMap* } m_globalTopic = getTopic(nullptr, "global/"); m_subscribeTopic = getTopic(nullptr, "#"); - m_mosquitto = nullptr; - if (mosquitto_lib_init() != MOSQ_ERR_SUCCESS) { - logOtherError("mqtt", "unable to initialize"); - } else { + if (check(mosquitto_lib_init(), "unable to initialize")) { signal(SIGPIPE, SIG_IGN); // needed before libmosquitto v. 1.1.3 ostringstream clientId; if (g_clientId) { @@ -394,8 +468,12 @@ MqttHandler::MqttHandler(UserInfo* userInfo, BusHandler* busHandler, MessageMap* } } if (m_mosquitto) { - /*mosquitto_log_init(m_mosquitto, MOSQ_LOG_DEBUG | MOSQ_LOG_ERR | MOSQ_LOG_WARNING - | MOSQ_LOG_NOTICE | MOSQ_LOG_INFO, MOSQ_LOG_STDERR);*/ +#if (LIBMOSQUITTO_VERSION_NUMBER >= 1003002) + check(mosquitto_threaded_set(m_mosquitto, true), "threaded_set"); +#endif +#if (LIBMOSQUITTO_VERSION_NUMBER >= 1004001) + check(mosquitto_opts_set(m_mosquitto, MOSQ_OPT_PROTOCOL_VERSION, (void*)(&g_version)), "opts_set protocol version"); +#endif if (g_username || g_password) { if (!g_username) { g_username = PACKAGE; @@ -423,6 +501,11 @@ MqttHandler::MqttHandler(UserInfo* userInfo, BusHandler* busHandler, MessageMap* logOtherError("mqtt", "unable to set TLS: %d", ret); } } +#endif +#if (LIBMOSQUITTO_VERSION_NUMBER >= 1003001) + if (g_logFromLib) { + mosquitto_log_callback_set(m_mosquitto, on_log); + } #endif mosquitto_connect_callback_set(m_mosquitto, on_connect); mosquitto_message_callback_set(m_mosquitto, on_message); @@ -436,16 +519,9 @@ MqttHandler::MqttHandler(UserInfo* userInfo, BusHandler* busHandler, MessageMap* logOtherError("mqtt", "unable to connect (invalid parameters)"); mosquitto_destroy(m_mosquitto); m_mosquitto = nullptr; - } else if (ret != MOSQ_ERR_SUCCESS) { + } else if (!check(ret, "unable to connect, retrying")) { m_connected = false; m_initialConnectFailed = g_ignoreInvalidParams; - 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 { m_connected = true; // assume success until connect_callback says otherwise logOtherDebug("mqtt", "connection requested"); @@ -473,7 +549,7 @@ void MqttHandler::notifyConnected() { const string sep = (g_publishFormat & OF_JSON) ? "\"" : ""; publishTopic(m_globalTopic+"version", sep + (PACKAGE_STRING "." REVISION) + sep, true); publishTopic(m_globalTopic+"running", "true", true); - mosquitto_subscribe(m_mosquitto, nullptr, m_subscribeTopic.c_str(), 0); + check(mosquitto_subscribe(m_mosquitto, nullptr, m_subscribeTopic.c_str(), 0), "subscribe"); } } @@ -646,45 +722,46 @@ void MqttHandler::run() { } void MqttHandler::handleTraffic(bool allowReconnect) { - if (m_mosquitto) { - int ret; + if (!m_mosquitto) { + return; + } + int ret; #if (LIBMOSQUITTO_MAJOR >= 1) - ret = mosquitto_loop(m_mosquitto, -1, 1); // waits up to 1 second for network traffic + ret = mosquitto_loop(m_mosquitto, -1, 1); // waits up to 1 second for network traffic #else - ret = mosquitto_loop(m_mosquitto, -1); // waits up to 1 second for network traffic + ret = mosquitto_loop(m_mosquitto, -1); // waits up to 1 second for network traffic #endif - if (!m_connected && ret == MOSQ_ERR_NO_CONN && allowReconnect) { - if (m_initialConnectFailed) { + if (!m_connected && ret == MOSQ_ERR_NO_CONN && allowReconnect) { + if (m_initialConnectFailed) { #if (LIBMOSQUITTO_MAJOR >= 1) - ret = mosquitto_connect(m_mosquitto, g_host, g_port, 60); + ret = mosquitto_connect(m_mosquitto, g_host, g_port, 60); #else - ret = mosquitto_connect(m_mosquitto, g_host, g_port, 60, true); + ret = mosquitto_connect(m_mosquitto, g_host, g_port, 60, true); #endif - if (ret == MOSQ_ERR_INVAL) { - logOtherError("mqtt", "unable to connect (invalid parameters), retrying"); - } - if (ret == MOSQ_ERR_SUCCESS) { - m_initialConnectFailed = false; - } - } else { - ret = mosquitto_reconnect(m_mosquitto); + if (ret == MOSQ_ERR_INVAL) { + logOtherError("mqtt", "unable to connect (invalid parameters), retrying"); + } + if (ret == MOSQ_ERR_SUCCESS) { + m_initialConnectFailed = false; } - } - if (!m_connected && ret == MOSQ_ERR_SUCCESS) { - m_connected = true; - logOtherNotice("mqtt", "connection re-established"); - } - if (!m_connected || ret == MOSQ_ERR_SUCCESS) { - return; - } - if (ret == MOSQ_ERR_NO_CONN || ret == MOSQ_ERR_CONN_LOST || ret == MOSQ_ERR_CONN_REFUSED) { - logOtherError("mqtt", "communication error: %s", ret == MOSQ_ERR_NO_CONN ? "not connected" - : (ret == MOSQ_ERR_CONN_LOST ? "connection lost" : "connection refused")); - m_connected = false; } else { - logOtherError("mqtt", "communication error: %d", ret); + ret = mosquitto_reconnect(m_mosquitto); } } + if (!m_connected && ret == MOSQ_ERR_SUCCESS) { + m_connected = true; + logOtherNotice("mqtt", "connection re-established"); + } + if (!m_connected || ret == MOSQ_ERR_SUCCESS) { + return; + } + if (ret == MOSQ_ERR_NO_CONN || ret == MOSQ_ERR_CONN_LOST || ret == MOSQ_ERR_CONN_REFUSED) { + logOtherError("mqtt", "communication error: %s", ret == MOSQ_ERR_NO_CONN ? "not connected" + : (ret == MOSQ_ERR_CONN_LOST ? "connection lost" : "connection refused")); + m_connected = false; + } else { + check(ret, "communication error"); + } } string MqttHandler::getTopic(const Message* message, const string& suffix, const string& fieldName) { @@ -745,9 +822,12 @@ void MqttHandler::publishMessage(const Message* message, ostringstream* updates) } void MqttHandler::publishTopic(const string& topic, const string& data, bool retain) { - logOtherDebug("mqtt", "publish %s %s", topic.c_str(), data.c_str()); - mosquitto_publish(m_mosquitto, nullptr, topic.c_str(), (uint32_t)data.size(), - reinterpret_cast(data.c_str()), 0, g_retain || retain); + const char* topicStr = topic.c_str(); + const char* dataStr = data.c_str(); + const size_t len = strlen(dataStr); + logOtherDebug("mqtt", "publish %s %s", topicStr, dataStr); + check(mosquitto_publish(m_mosquitto, nullptr, topic.c_str(), (uint32_t)len, + reinterpret_cast(dataStr), 0, g_retain || retain), "publish"); } } // namespace ebusd