add qos param, add global topic param and allow reduction to alive only (addresses #355), allow disabling defaults for topic (solves #355), optimized

This commit is contained in:
John
2022-02-05 13:05:23 +01:00
parent 6f56ca6331
commit 1d232e211f
2 changed files with 156 additions and 47 deletions
+134 -45
View File
@@ -36,8 +36,10 @@ using std::dec;
#define O_USER (O_CLID+1)
#define O_PASS (O_USER+1)
#define O_TOPI (O_PASS+1)
#define O_RETA (O_TOPI+1)
#define O_INTF (O_RETA+1)
#define O_GTOP (O_TOPI+1)
#define O_RETA (O_GTOP+1)
#define O_PQOS (O_RETA+1)
#define O_INTF (O_PQOS+1)
#define O_JSON (O_INTF+1)
#define O_LOGL (O_JSON+1)
#define O_VERS (O_LOGL+1)
@@ -61,7 +63,9 @@ static const struct argp_option g_mqtt_argp_options[] = {
{"mqttpass", O_PASS, "PASSWORD", 0, "Use PASSWORD when connecting to MQTT broker (no default)", 0 },
{"mqtttopic", O_TOPI, "TOPIC", 0,
"Use MQTT TOPIC (prefix before /%circuit/%name or complete format) [ebusd]", 0 },
{"mqttglobal", O_GTOP, "TOPIC", 0, "Use TOPIC for global data (default is \"global/\" suffix to mqtttopic prefix)", 0 },
{"mqttretain", O_RETA, nullptr, 0, "Retain all topics instead of only selected global ones", 0 },
{"mqttqos", O_PQOS, "QOS", 0, "Set the QoS value for all topics (0-2) [0]", 0 },
{"mqttint", O_INTF, "FILE", 0, "Read MQTT integration settings from FILE (no default)", 0 },
{"mqttjson", O_JSON, nullptr, 0, "Publish in JSON format instead of strings", 0 },
{"mqttverbose", O_VERB, nullptr, 0, "Publish all available attributes", 0 },
@@ -92,8 +96,10 @@ static const char* g_clientId = nullptr; //!< optional clientid override for MQ
static const char* g_username = nullptr; //!< optional user name for MQTT broker (no default)
static const char* g_password = nullptr; //!< optional password for MQTT broker (no default)
static const char* g_topic = nullptr; //!< optional topic template
static const char* g_globalTopic = nullptr; //!< optional global topic
static const char* g_integrationFile = nullptr; //!< the integration settings file
static bool g_retain = false; //!< whether to retail all topics
static int g_qos = 0; //!< the qos value for all topics
static OutputFormat g_publishFormat = OF_NONE; //!< the OutputFormat for publishing messages
#if (LIBMOSQUITTO_VERSION_NUMBER >= 1003001)
static bool g_logFromLib = false; //!< log library events
@@ -173,9 +179,15 @@ static error_t mqtt_parse_opt(int key, char *arg, struct argp_state *state) {
break;
case O_TOPI: // --mqtttopic=ebusd
if (arg == nullptr || arg[0] == 0 || strchr(arg, '#') || strchr(arg, '+') || arg[strlen(arg)-1] == '/') {
if (arg == nullptr || arg[0] == 0 || strchr(arg, '+') || arg[strlen(arg)-1] == '/') {
argp_error(state, "invalid mqtttopic");
return EINVAL;
} else {
char *pos = strchr(arg, '#');
if (pos && (pos == arg || pos[1])) { // allow # only at very last position (to indicate not using any default)
argp_error(state, "invalid mqtttopic");
return EINVAL;
}
}
if (g_topic) {
argp_error(state, "duplicate mqtttopic");
@@ -190,10 +202,26 @@ static error_t mqtt_parse_opt(int key, char *arg, struct argp_state *state) {
}
break;
case O_GTOP: // --mqttglobal=global/
if (arg == nullptr || strchr(arg, '+') || strchr(arg, '#')) {
argp_error(state, "invalid mqttglobal");
return EINVAL;
}
g_globalTopic = arg;
break;
case O_RETA: // --mqttretain
g_retain = true;
break;
case O_PQOS: // --mqttqos=0
g_qos = parseSignedInt(arg, 10, 0, 2, &result);
if (result != RESULT_OK) {
argp_error(state, "invalid mqttqos value");
return EINVAL;
}
break;
case O_INTF: // --mqttint=/etc/ebusd/mqttint.cfg
if (arg == nullptr || arg[0] == 0 || strcmp("/", arg) == 0) {
argp_error(state, "invalid mqttint file");
@@ -501,6 +529,26 @@ string MqttReplacer::get(const map<string, string>& values, bool untilFirstEmpty
return str;
}
string MqttReplacer::get(const string& circuit, const string& name, const string& fieldName) const {
map <string, string> values;
values["circuit"] = circuit;
values["name"] = name;
if (!fieldName.empty()) {
values["field"] = fieldName;
}
return get(values, true);
}
string MqttReplacer::get(const Message* message, const string& fieldName) const {
map<string, string> values;
values["circuit"] = message->getCircuit();
values["name"] = message->getName();
if (!fieldName.empty()) {
values["field"] = fieldName;
}
return get(message->getCircuit(), message->getName(), fieldName);
}
bool MqttReplacer::isReducable(const map<string, string>& values) const {
for (const auto &it: m_parts) {
if (it.second < 0) {
@@ -890,23 +938,32 @@ MqttHandler::MqttHandler(UserInfo* userInfo, BusHandler* busHandler, MessageMap*
// determine topic and prefix
MqttReplacer& topic = m_replacers.get("topic");
if (g_topic) {
string str = g_topic;
bool noDefault = str[str.size()-1]=='#';
if (noDefault) {
str.resize(str.size()-1);
}
if (hasIntegration && !topic.empty()) {
// topic defined in cmdline and integration file.
if (strchr(g_topic, '%')) {
if (str.find('%')!=string::npos) {
// cmdline topic is more than just a prefix => override integration topic completely
topic.parse(g_topic);
topic.parse(str);
} else {
// cmdline topic is only the prefix, use it
string prefix = g_topic;
m_replacers.set("prefix", prefix);
m_replacers.set("prefixn", removeTrailingNonTopicPart(prefix));
m_replacers.set("prefix", str);
m_replacers.set("prefixn", removeTrailingNonTopicPart(str));
}
} else {
topic.parse(g_topic);
topic.parse(str);
}
if (!noDefault) {
topic.ensureDefault();
}
} else {
topic.ensureDefault();
}
topic.ensureDefault();
m_publishByField = topic.has("field");
m_staticTopic = !topic.has("name");
m_publishByField = !m_staticTopic && topic.has("field");
if (hasIntegration) {
m_replacers.set("version", PACKAGE_VERSION);
if (m_replacers["prefix"].empty()) {
@@ -956,7 +1013,17 @@ MqttHandler::MqttHandler(UserInfo* userInfo, BusHandler* busHandler, MessageMap*
m_hasDefinitionFieldsPayload = m_replacers.uses("fields_payload");
m_subscribeConfigRestartTopic = m_replacers.get("config_restart-topic", false, false);
m_subscribeConfigRestartPayload = m_replacers.get("config_restart-payload", false, false);
m_globalTopic = getTopic(nullptr, "global/");
if (g_globalTopic) {
m_globalTopic.parse(g_globalTopic, true);
} else {
m_globalTopic.parse(getTopic(nullptr, "%circuit/%name"), true);
}
if (m_globalTopic.has("circuit")) {
map<string, string> values;
values["circuit"] = "global";
string result;
m_globalTopic.reduce(values, result);
}
m_subscribeTopic = getTopic(nullptr, "#");
if (check(mosquitto_lib_init(), "unable to initialize")) {
signal(SIGPIPE, SIG_IGN); // needed before libmosquitto v. 1.1.3
@@ -989,7 +1056,7 @@ MqttHandler::MqttHandler(UserInfo* userInfo, BusHandler* busHandler, MessageMap*
logOtherError("mqtt", "unable to set username/password, trying without");
}
}
string willTopic = m_globalTopic+"running";
string willTopic = m_globalTopic.get("", "running");
string willData = "false";
size_t len = willData.length();
#if (LIBMOSQUITTO_MAJOR >= 1)
@@ -1059,11 +1126,15 @@ void MqttHandler::start() {
void MqttHandler::notifyConnected() {
if (m_mosquitto && isRunning()) {
const string sep = (g_publishFormat & OF_JSON) ? "\"" : "";
publishTopic(m_globalTopic+"version", sep + (PACKAGE_STRING "." REVISION) + sep, true);
publishTopic(m_globalTopic+"running", "true", true);
check(mosquitto_subscribe(m_mosquitto, nullptr, m_subscribeTopic.c_str(), 0), "subscribe");
if (!m_subscribeConfigRestartTopic.empty()) {
check(mosquitto_subscribe(m_mosquitto, nullptr, m_subscribeConfigRestartTopic.c_str(), 0), "subscribe definition");
if (m_globalTopic.has("name")) {
publishTopic(m_globalTopic.get("", "version"), sep + (PACKAGE_STRING "." REVISION) + sep, true);
}
publishTopic(m_globalTopic.get("", "running"), "true", true);
if (!m_staticTopic) {
check(mosquitto_subscribe(m_mosquitto, nullptr, m_subscribeTopic.c_str(), 0), "subscribe");
if (!m_subscribeConfigRestartTopic.empty()) {
check(mosquitto_subscribe(m_mosquitto, nullptr, m_subscribeConfigRestartTopic.c_str(), 0), "subscribe definition");
}
}
}
}
@@ -1174,16 +1245,20 @@ void MqttHandler::notifyTopic(const string& topic, const string& data) {
void MqttHandler::notifyUpdateCheckResult(const string& checkResult) {
if (checkResult != m_lastUpdateCheckResult) {
m_lastUpdateCheckResult = checkResult;
const string sep = (g_publishFormat & OF_JSON) ? "\"" : "";
publishTopic(m_globalTopic+"updatecheck", sep + (checkResult.empty() ? "OK" : checkResult) + sep, true);
if (m_globalTopic.has("name")) {
const string sep = (g_publishFormat & OF_JSON) ? "\"" : "";
publishTopic(m_globalTopic.get("", "updatecheck"), sep + (checkResult.empty() ? "OK" : checkResult) + sep, true);
}
}
}
void MqttHandler::notifyScanStatus(const string& scanStatus) {
if (scanStatus != m_lastScanStatus) {
m_lastScanStatus = scanStatus;
const string sep = (g_publishFormat & OF_JSON) ? "\"" : "";
publishTopic(m_globalTopic+"scan", sep + (scanStatus.empty() ? "OK" : scanStatus) + sep, true);
if (m_globalTopic.has("name")) {
const string sep = (g_publishFormat & OF_JSON) ? "\"" : "";
publishTopic(m_globalTopic.get("", "scan"), sep + (scanStatus.empty() ? "OK" : scanStatus) + sep, true);
}
}
}
@@ -1204,8 +1279,9 @@ void splitFields(const string& str, vector<string>* row) {
void MqttHandler::run() {
time_t lastTaskRun, now, start, lastSignal = 0, lastUpdates = 0;
bool signal = false;
string signalTopic = m_globalTopic + "signal";
string uptimeTopic = m_globalTopic + "uptime";
bool globalHasName = m_globalTopic.has("name");
string signalTopic = m_globalTopic.get("", "signal");
string uptimeTopic = m_globalTopic.get("", "uptime");
ostringstream updates;
unsigned int filterPriority = 0;
bool filterSeen = false;
@@ -1255,16 +1331,23 @@ void MqttHandler::run() {
time_t uptime = now - start;
updates.str("");
updates.clear();
updates << dec << static_cast<unsigned>(uptime);
publishTopic(uptimeTopic, updates.str());
if (globalHasName) {
updates << dec << static_cast<unsigned>(uptime);
publishTopic(uptimeTopic, updates.str());
}
}
if (m_connected && m_definitionsSince == 0) {
publishDefinition(m_replacers, "def_global_running-", m_globalTopic + "running", "global", "running", "def_global-");
publishDefinition(m_replacers, "def_global_version-", m_globalTopic + "version", "global", "version", "def_global-");
publishDefinition(m_replacers, "def_global_signal-", signalTopic, "global", "signal", "def_global-");
publishDefinition(m_replacers, "def_global_uptime-", uptimeTopic, "global", "uptime", "def_global-");
publishDefinition(m_replacers, "def_global_updatecheck-", m_globalTopic + "updatecheck", "global", "updatecheck", "def_global-");
publishDefinition(m_replacers, "def_global_scan-", m_globalTopic + "scan", "global", "scan", "def_global-");
publishDefinition(m_replacers, "def_global_running-", m_globalTopic.get("", "running"), "global", "running", "def_global-");
if (globalHasName) {
publishDefinition(m_replacers, "def_global_version-", m_globalTopic.get("", "version"), "global", "version",
"def_global-");
publishDefinition(m_replacers, "def_global_signal-", signalTopic, "global", "signal", "def_global-");
publishDefinition(m_replacers, "def_global_uptime-", uptimeTopic, "global", "uptime", "def_global-");
publishDefinition(m_replacers, "def_global_updatecheck-", m_globalTopic.get("", "updatecheck"), "global",
"updatecheck", "def_global-");
publishDefinition(m_replacers, "def_global_scan-", m_globalTopic.get("", "scan"), "global", "scan",
"def_global-");
}
m_definitionsSince = 1;
}
if (m_connected && m_hasDefinitionTopic) {
@@ -1446,12 +1529,16 @@ void MqttHandler::run() {
lastSignal = now;
if (!signal || reconnected) {
signal = true;
publishTopic(signalTopic, "true", true);
if (globalHasName) {
publishTopic(signalTopic, "true", true);
}
}
} else {
if (signal || reconnected) {
signal = false;
publishTopic(signalTopic, "false", true);
if (globalHasName) {
publishTopic(signalTopic, "false", true);
}
}
}
}
@@ -1483,8 +1570,10 @@ void MqttHandler::run() {
break;
}
}
publishTopic(signalTopic, "false", true);
publishTopic(m_globalTopic+"scan", "", true); // clear retain of scan status
if (globalHasName) {
publishTopic(signalTopic, "false", true);
publishTopic(m_globalTopic.get("", "scan"), "", true); // clear retain of scan status
}
}
void MqttHandler::publishDefinition(MqttReplacers values, const string& prefix, const string& topic,
@@ -1584,16 +1673,10 @@ bool MqttHandler::handleTraffic(bool allowReconnect) {
}
string MqttHandler::getTopic(const Message* message, const string& suffix, const string& fieldName) {
if (!message) {
if (!message || m_staticTopic) {
return m_replacers.get("topic", true) + suffix;
}
map <string, string> values;
values["circuit"] = message->getCircuit();
values["name"] = message->getName();
if (!fieldName.empty()) {
values["field"] = fieldName;
}
return m_replacers.get("topic").get(values, true) + suffix;
return m_replacers.get("topic").get(message, fieldName) + suffix;
}
void MqttHandler::publishMessage(const Message* message, ostringstream* updates, bool includeWithoutData) {
@@ -1607,6 +1690,9 @@ void MqttHandler::publishMessage(const Message* message, ostringstream* updates,
}
if (json) {
*updates << "{";
if (m_staticTopic) {
*updates << "\"circuit\":\"" << message->getCircuit() << "\",\"name\":\"" << message->getName() << "\",\"fields\":{";
}
}
result_t result = message->decodeLastData(false, nullptr, -1, outputFormat, updates);
if (result != RESULT_OK) {
@@ -1615,6 +1701,9 @@ void MqttHandler::publishMessage(const Message* message, ostringstream* updates,
return;
}
if (json) {
if (m_staticTopic) {
*updates << "}";
}
*updates << "}";
}
publishTopic(getTopic(message), updates->str());
@@ -1647,7 +1736,7 @@ void MqttHandler::publishTopic(const string& topic, const string& data, bool ret
const size_t len = strlen(dataStr);
logOtherDebug("mqtt", "publish %s %s", topicStr, dataStr);
check(mosquitto_publish(m_mosquitto, nullptr, topicStr, (uint32_t)len,
reinterpret_cast<const uint8_t*>(dataStr), 0, g_retain || retain), "publish");
reinterpret_cast<const uint8_t*>(dataStr), g_qos, g_retain || retain), "publish");
}
void MqttHandler::publishEmptyTopic(const string& topic) {
+22 -2
View File
@@ -111,6 +111,23 @@ class MqttReplacer {
*/
string get(const map<string, string>& values, bool untilFirstEmpty = true, bool onlyAlphanum = false) const;
/**
* Get the replaced template string.
* @param circuit the circuit name for replacement.
* @param name the message name for replacement.
* @param fieldName the field name for replacement.
* @return the replaced template string.
*/
string get(const string& circuit, const string& name, const string& fieldName="") const;
/**
* Get the replaced template string.
* @param message the Message from which to extract the values for replacement.
* @param fieldName the field name for replacement.
* @return the replaced template string.
*/
string get(const Message* message, const string& fieldName="") const;
/**
* Check if the fields can be reduced to a constant value.
* @param values the named values for replacement.
@@ -342,12 +359,15 @@ class MqttHandler : public DataSink, public DataSource, public WaitThread {
/** the @a MessageMap instance. */
MessageMap* m_messages;
/** the global topic prefix. */
string m_globalTopic;
/** the global topic replacer. */
MqttReplacer m_globalTopic;
/** the topic to subscribe to. */
string m_subscribeTopic;
/** whether to use a single topic for all messages. */
bool m_staticTopic;
/** whether to publish a separate topic for each message field. */
bool m_publishByField;