added --mqttlog and --mqttversion options, check return code of calls to mosquitto library, log library version, set threaded mode
This commit is contained in:
+129
-49
@@ -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<string> g_topicStrs;
|
||||
static vector<string> 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<DataHandler*>* 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<const uint8_t*>(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<const uint8_t*>(dataStr), 0, g_retain || retain), "publish");
|
||||
}
|
||||
|
||||
} // namespace ebusd
|
||||
|
||||
Reference in New Issue
Block a user