diff --git a/src/ebusd/mqtthandler.cpp b/src/ebusd/mqtthandler.cpp index ab5dcb79..a8fbb855 100755 --- a/src/ebusd/mqtthandler.cpp +++ b/src/ebusd/mqtthandler.cpp @@ -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& values, bool untilFirstEmpty return str; } +string MqttReplacer::get(const string& circuit, const string& name, const string& fieldName) const { + map 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 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& 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 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* 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(uptime); - publishTopic(uptimeTopic, updates.str()); + if (globalHasName) { + updates << dec << static_cast(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 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(dataStr), 0, g_retain || retain), "publish"); + reinterpret_cast(dataStr), g_qos, g_retain || retain), "publish"); } void MqttHandler::publishEmptyTopic(const string& topic) { diff --git a/src/ebusd/mqtthandler.h b/src/ebusd/mqtthandler.h index bedbb076..c77eaf35 100644 --- a/src/ebusd/mqtthandler.h +++ b/src/ebusd/mqtthandler.h @@ -111,6 +111,23 @@ class MqttReplacer { */ string get(const map& 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;