diff --git a/src/ebusd/mqtthandler.cpp b/src/ebusd/mqtthandler.cpp index 39970563..8a3bbfc2 100755 --- a/src/ebusd/mqtthandler.cpp +++ b/src/ebusd/mqtthandler.cpp @@ -36,7 +36,8 @@ using std::dec; #define O_PASS (O_USER+1) #define O_TOPI (O_PASS+1) #define O_RETA (O_TOPI+1) -#define O_JSON (O_RETA+1) +#define O_INTF (O_RETA+1) +#define O_JSON (O_INTF+1) #define O_LOGL (O_JSON+1) #define O_VERS (O_LOGL+1) #define O_IGIN (O_VERS+1) @@ -53,13 +54,14 @@ static const struct argp_option g_mqtt_argp_options[] = { {nullptr, 0, nullptr, 0, "MQTT options:", 1 }, {"mqtthost", O_HOST, "HOST", 0, "Connect to MQTT broker on HOST [localhost]", 0 }, {"mqttport", O_PORT, "PORT", 0, "Connect to MQTT broker on PORT (usually 1883), 0 to disable [0]", 0 }, - {"mqttclientid", O_CLID, "ID", 0, "Set client ID for connection to MQTT broker [" PACKAGE_NAME "_" + {"mqttclientid", O_CLID, "ID", 0, "Set client ID for connection to MQTT broker [" PACKAGE_NAME "_" PACKAGE_VERSION "_]", 0 }, {"mqttuser", O_USER, "USER", 0, "Connect as USER to MQTT broker (no default)", 0 }, {"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 }, {"mqttretain", O_RETA, nullptr, 0, "Retain all topics instead of only selected global ones", 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 }, #if (LIBMOSQUITTO_VERSION_NUMBER >= 1003001) @@ -88,10 +90,8 @@ static uint16_t g_port = 0; //!< optional port of MQTT broker, 0 t static const char* g_clientId = nullptr; //!< optional clientid override for MQTT broker 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) -/** the MQTT topic string parts. */ -static vector g_topicStrs; -/** the MQTT topic field parts. */ -static vector g_topicFields; +static MqttReplacer* g_topicReplacer = nullptr; //!< the topic replacer +static const char* g_integrationFile = nullptr; //!< the integration settings file static bool g_retain = false; //!< whether to retail all topics static OutputFormat g_publishFormat = OF_NONE; //!< the OutputFormat for publishing messages #if (LIBMOSQUITTO_VERSION_NUMBER >= 1003001) @@ -112,8 +112,6 @@ static const char* g_keypass = nullptr; //!< client key file password for TLS static bool g_insecure = false; //!< whether to allow insecure TLS connection #endif -bool parseTopic(const string& topic, vector* strs, vector* fields); - static char* replaceSecret(char *arg) { char* ret = strdup(arg); int cnt = 0; @@ -178,7 +176,12 @@ static error_t mqtt_parse_opt(int key, char *arg, struct argp_state *state) { argp_error(state, "invalid mqtttopic"); return EINVAL; } - if (!parseTopic(arg, &g_topicStrs, &g_topicFields)) { + if (g_topicReplacer) { + argp_error(state, "duplicate mqtttopic"); + return EINVAL; + } + g_topicReplacer = MqttReplacer::create(arg); + if (!g_topicReplacer) { argp_error(state, "malformed mqtttopic"); } break; @@ -187,6 +190,14 @@ static error_t mqtt_parse_opt(int key, char *arg, struct argp_state *state) { g_retain = true; break; + case O_INTF: // --mqttint=/etc/ebusd/mqttint.cfg + if (arg == nullptr || arg[0] == 0 || strcmp("/", arg) == 0) { + argp_error(state, "invalid mqttint file"); + return EINVAL; + } + g_integrationFile = arg; + break; + case O_JSON: // --mqttjson g_publishFormat |= OF_JSON|OF_NAMES; break; @@ -274,9 +285,6 @@ static const struct argp_child g_mqtt_argp_child = {&g_mqtt_argp, 0, "", 1}; const struct argp_child* mqtthandler_getargs() { - if (g_topicStrs.empty()) { - g_topicStrs.push_back(PACKAGE); - } return &g_mqtt_argp_child; } @@ -326,50 +334,312 @@ static const char* knownFieldNames[] = { /** the number of known field names. */ static const size_t knownFieldCount = sizeof(knownFieldNames) / sizeof(char*); - -/** - * Parse the topic template. - * @param topic the topic template. - * @param strs the @a vector to which the string parts shall be added. - * @param fields the @a vector to which the field parts shall be added. - * @return true on success, false on malformed topic template. - */ -bool parseTopic(const string& topic, vector* strs, vector* fields) { - size_t lastpos = 0; - size_t end = topic.length(); - vector columns; - strs->clear(); - fields->clear(); - for (size_t pos=topic.find('%', lastpos); pos != string::npos; ) { - size_t idx = knownFieldCount; - size_t len = 0; - for (size_t i = 0; i < knownFieldCount; i++) { - len = strlen(knownFieldNames[i]); - if (topic.substr(pos+1, len) == knownFieldNames[i]) { - idx = i; - break; - } - } - if (idx == knownFieldCount) { // TODO could allow custom attributes here - return false; - } - string fieldName = knownFieldNames[idx]; - for (const auto& it : *fields) { - if (it == fieldName) { - return false; // duplicate column - } - } - strs->push_back(topic.substr(lastpos, pos-lastpos)); - fields->push_back(fieldName); - lastpos = pos+1+len; - pos = topic.find('%', lastpos); +std::pair makeField(string name, bool isField) { + if (!isField) { + return {name, -1}; } - if (lastpos < end) { - strs->push_back(topic.substr(lastpos, end-lastpos)); + for (int idx = 0; idx < (int)knownFieldCount; idx++) { + if (name==knownFieldNames[idx]) { + return {name, idx}; + } + } + return {name, knownFieldCount}; +} + + +MqttReplacer::MqttReplacer(bool fillDefault) { + if (fillDefault) { + ensureDefault(); + } +} + +bool MqttReplacer::parse(const string& templateStr, bool onlyKnown, bool noKnownDuplicates) { + m_parts.clear(); + size_t end = templateStr.length(); + bool inField = false; + ostringstream stack; + for (size_t pos = 0; pos < end+1; pos++) { + char ch = pos= 'a' && ch <= 'z') || (ch >= 'A' && ch <= 'Z') || ch == '_')) { + // invalid field character + if (stack.tellp()>0) { + m_parts.push_back(makeField(stack.str(), true)); + stack.str(""); + } + inField = false; + } + stack << ch; + } + } + if (onlyKnown || noKnownDuplicates) { + int foundMask = 0; + int knownCount = knownFieldCount; + for (const auto &it : m_parts) { + if (it.second < 0) { + continue; // unknown field + } + if (onlyKnown && it.second >= knownCount) { + return false; + } + if (noKnownDuplicates && it.second < knownCount) { + int bit = 1<parse(templateStr, onlyKnown, noKnownDuplicates)) { + if (ensureDefault) { + ret->ensureDefault(); + } + return ret; + } + delete ret; + return nullptr; +} + +void MqttReplacer::normalize(string& str) { + transform(str.begin(), str.end(), str.begin(), [](unsigned char c){ + return isalnum(c) ? c : '_'; + }); +} + +void MqttReplacer::ensureDefault() { + if (m_parts.empty()) { + m_parts.emplace_back(string(PACKAGE) + "/", -1); + } else if (m_parts.size()==1 && m_parts[0].second<0 && m_parts[0].first.find('/')==string::npos) { + m_parts[0] = {m_parts[0].first + "/", -1}; // ensure trailing slash + } + if (!has("circuit")) { + m_parts.emplace_back("circuit", 0); // index of circuit in knownFieldNames + m_parts.emplace_back("/", -1); + } + if (!has("name")) { + m_parts.emplace_back("name", 1); // index of name in knownFieldNames + } +} + +bool MqttReplacer::has(const string& field) const { + for (const auto &it: m_parts) { + if (it.second >= 0 && it.first == field) { + return true; + } + } + return false; +} + +string MqttReplacer::get(const map& values, bool untilFirstEmpty, bool onlyAlphanum) const { + ostringstream ret; + for (const auto &it: m_parts) { + if (it.second < 0) { + ret << it.first; + continue; + } + const auto pos = values.find(it.first); + if (pos==values.cend()) { + if (untilFirstEmpty) { + break; + } + } else if (untilFirstEmpty && pos->second.empty()) { + break; + } else { + ret << pos->second; + } + } + if (!onlyAlphanum) { + return ret.str(); + } + string str = ret.str(); + normalize(str); + return str; +} + +bool MqttReplacer::isReducable(const map& values) const { + for (const auto &it: m_parts) { + if (it.second < 0) { + continue; + } + const auto pos = values.find(it.first); + if (pos==values.cend()) { + return false; + } + } + return true; +} + +bool MqttReplacer::reduce(const map& values, string& result, bool onlyAlphanum) const { + ostringstream ret; + for (const auto &it: m_parts) { + if (it.second < 0) { + ret << it.first; + continue; + } + const auto pos = values.find(it.first); + if (pos==values.cend()) { + result = ret.str(); + return false; + } + ret << pos->second; + } + result = ret.str(); + if (onlyAlphanum) { + normalize(result); + } + return true; +} + +ssize_t MqttReplacer::matchTopic(const string& remain, string* circuit, string* name, string* field) const { + size_t last = 0; + size_t count = m_parts.size(); + size_t idx; + for (idx = 0; idx < count; idx++) { + const auto part = m_parts[idx]; + if (part.second < 0) { + if (remain.substr(last, part.first.length()) != part.first) { + return static_cast(idx); + } + last += part.first.length(); + continue; + } + string value; + if (idx+1 < count) { + // todo require topic fields to be separated by non-empty string? e.g. %circuit%name is not parseable here + size_t pos = remain.find(m_parts[idx+1].first, last); + if (pos == string::npos) { + // next part not found + return -static_cast(idx)-1; + } + value = remain.substr(last, pos-last); + } else { + // last part is a field name + if (remain.find('/', last) != string::npos) { + // non-name in remainder found + return -static_cast(idx)-1; + } + value = remain; + } + last += value.length(); + switch (part.second) { + case 0: *circuit = value; break; + case 1: *name = value; break; + case 2: *field = value; break; + default: // unknown field + break; + } + } + return static_cast(idx); +} + +static const string EMPTY = ""; + +const string& MqttReplacers::operator[](const string& key) const { + auto itc = m_constants.find(key); + if (itc==m_constants.end()) { + return EMPTY; + } + return itc->second; +} + +bool MqttReplacers::uses(const string& field) const { + for (const auto &it : m_replacers) { + if (it.second.has(field)) { + return true; + } + } + return false; +} + +MqttReplacer& MqttReplacers::get(const string& key) { + return m_replacers[key]; +} + +string MqttReplacers::get(const string& key, bool untilFirstEmpty, bool onlyAlphanum) const { + auto itc = m_constants.find(key); + if (itc!=m_constants.end()) { + return itc->second; + } + auto itv = m_replacers.find(key); + if (itv==m_replacers.end()) { + return ""; + } + return itv->second.get(m_constants, untilFirstEmpty, onlyAlphanum); +} + +bool MqttReplacers::set(const string& key, const string& value, bool removeReplacer) { + m_constants[key] = value; + if (removeReplacer) { + m_replacers.erase(key); + } + if (key.find('_')!=string::npos) { + return false; + } + string upper = key; + transform(upper.begin(), upper.end(), upper.begin(), ::toupper); + if (upper==key) { + return false; + } + string val = value; + MqttReplacer::normalize(val); + m_constants[upper] = val; + if (removeReplacer) { + m_replacers.erase(upper); + } + return true; +} + +void MqttReplacers::set(const string& key, int value) { + std::ostringstream str; + str << static_cast(value); + m_constants[key] = str.str(); +} + +void MqttReplacers::reduce() { + // iterate through variables and reduce as many to constants as possible + bool reduced = false; + do { + reduced = false; + for (auto it = m_replacers.begin(); it != m_replacers.end(); ) { + string str; + if (!it->second.isReducable(m_constants) + || !it->second.reduce(m_constants, str)) { + ++it; + continue; + } + bool restart = set(it->first, str, false); + it = m_replacers.erase(it); + reduced = true; + if (restart) { + string upper = it->first; + transform(upper.begin(), upper.end(), upper.begin(), ::toupper); + if (m_replacers.erase(upper)>0) { + break; // restart as iterator is now invalid + } + } + } + } while (reduced); +} + + #if (LIBMOSQUITTO_MAJOR >= 1) int on_keypassword(char *buf, int size, int rwflag, void *userdata) { if (!g_keypass) { @@ -444,32 +714,75 @@ void on_message( handler->notifyTopic(topic, data); } +void MqttHandler::parseIntegration(const string& line) { + if (line.empty()) { + return; + } + size_t pos = line.find('='); + if (pos==string::npos) { + return; + } + string key = line.substr(0, pos); + FileReader::trim(&key); + string value = line.substr(pos+1); + FileReader::trim(&value); + if (value.find('%')==string::npos) { + m_replacers.set(key, value); // constant value + } else { + // simple variable + m_replacers.get(key).parse(value, false, false); + } +} MqttHandler::MqttHandler(UserInfo* userInfo, BusHandler* busHandler, MessageMap* messages) : DataSink(userInfo, "mqtt"), DataSource(busHandler), WaitThread(), m_messages(messages), m_connected(false), m_initialConnectFailed(false), m_lastUpdateCheckResult("."), m_lastScanStatus("."), m_lastErrorLogTime(0) { - m_publishByField = false; + m_definitionsSince = 0; m_mosquitto = nullptr; - if (g_topicFields.empty()) { - if (g_topicStrs.empty()) { - g_topicStrs.push_back(""); + if (!g_topicReplacer) { + g_topicReplacer = new MqttReplacer(); + } + m_publishByField = g_topicReplacer->has("field"); + m_replacers.get("mqtttopic") = *g_topicReplacer; + if (g_integrationFile != nullptr) { + std::ifstream stream; + stream.open(g_integrationFile, std::ifstream::in); + if (!stream.is_open()) { + logOtherError("mqtt", "unable to open integration file %s", g_integrationFile); } else { - string str = g_topicStrs[0]; - if (str.empty() || str[str.length()-1] != '/') { - g_topicStrs[0] = str+"/"; - } - } - g_topicFields.push_back("circuit"); - g_topicStrs.push_back("/"); - g_topicFields.push_back("name"); - } else { - for (size_t i = 0; i < g_topicFields.size(); i++) { - if (g_topicFields[i] == "field") { - m_publishByField = true; - break; + m_replacers.set("version", PACKAGE_VERSION); + string line = m_replacers.get("mqtttopic", true); + m_replacers.set("prefix", line); + size_t pos = line.find_last_not_of("/_"); + m_replacers.set("prefixn", line.substr(0, pos+1)); + string last; + while (stream.peek() != EOF && getline(stream, line)) { + if (line.empty()) { + parseIntegration(last); + last = ""; + continue; + } + if (line[0] == '#') { + // only ignore it to allow commented lines in the middle of e.g. payload + continue; + } + if (last.empty()) { + last = line; + } else if (line[0] == '\t' || line[0] == ' ') { // continuation + last += "\n" + line; + } else { + parseIntegration(last); + last = line; + } } + parseIntegration(last); + m_replacers.reduce(); } } + m_hasDefinitionTopic = !m_replacers.get("definition_topic", false, false).empty(); + 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/"); m_subscribeTopic = getTopic(nullptr, "#"); if (check(mosquitto_lib_init(), "unable to initialize")) { @@ -576,6 +889,9 @@ void MqttHandler::notifyConnected() { 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"); + } } } @@ -584,6 +900,12 @@ void MqttHandler::notifyTopic(const string& topic, const string& data) { if (pos == string::npos) { return; } + if (!m_subscribeConfigRestartTopic.empty() && topic == m_subscribeConfigRestartTopic) { + if (m_subscribeConfigRestartPayload.empty() || data == m_subscribeConfigRestartPayload) { + m_definitionsSince = 0; + } + return; + } string direction = topic.substr(pos+1); if (direction.empty()) { return; @@ -595,60 +917,10 @@ void MqttHandler::notifyTopic(const string& topic, const string& data) { } logOtherDebug("mqtt", "received topic %s with data %s", topic.c_str(), data.c_str()); - string remain = topic.substr(0, pos); - size_t last = 0; - string circuit, name; - bool finalField = false; - for (size_t idx = 0; idx < g_topicStrs.size()+1 && !finalField; idx++) { - string field; - string chk; - if (idx < g_topicStrs.size()) { - chk = g_topicStrs[idx]; - pos = remain.find(chk, last); - if (pos == string::npos) { - if (!isList) { - return; - } - if (idx == 0 && remain+"/" == chk) { // check for only first prefix, e.g. "ebusd/" - break; - } - pos = remain.size(); - finalField = true; - } - } else if (idx-1 < g_topicFields.size()) { - pos = remain.size(); - } else if (last < remain.size()) { - if (!isList) { - return; - } - break; - } else { - break; - } - field = remain.substr(last, pos-last); - last = pos+chk.size(); - if (idx == 0) { - if (pos > 0) { - return; - } - } else { - if (field.empty()) { - if (!isList) { - return; - } - continue; - } - string fieldName = g_topicFields[idx-1]; - if (fieldName == "circuit") { - circuit = field; - } else if (fieldName == "name") { - name = field; - } else if (fieldName == "field") { - // field = field; // TODO add support for writing a single field - } else { - return; - } - } + string circuit, name, field; + ssize_t match = g_topicReplacer->matchTopic(topic.substr(0, pos), &circuit, &name, &field); + if (match<0 && !isList) { // TODO + logOtherError("mqtt", "received unmatchable topic %s", topic.c_str()); } if (isList) { logOtherInfo("mqtt", "received list topic for %s %s", circuit.c_str(), name.c_str()); @@ -664,7 +936,7 @@ void MqttHandler::notifyTopic(const string& topic, const string& data) { m_messages->findAll(circuit, name, m_levels, !(circuitPrefix || namePrefix), true, true, true, true, true, 0, 0, false, &messages); bool onlyWithData = !data.empty(); - for (const auto message : messages) { + for (const auto& message : messages) { if ((circuitPrefix && ( message->getCircuit().substr(0, circuit.length()) != circuit || (!namePrefix && name.length() > 0 && message->getName() != name))) @@ -767,12 +1039,100 @@ void MqttHandler::run() { lastTaskRun = now; } else if (now > lastTaskRun+15) { allowReconnect = true; - sendSignal = true; - time_t uptime = now-start; - updates.str(""); - updates.clear(); - updates << dec << static_cast(uptime); - publishTopic(uptimeTopic, updates.str()); + if (m_connected) { + sendSignal = true; + time_t uptime = now - start; + updates.str(""); + updates.clear(); + updates << dec << static_cast(uptime); + publishTopic(uptimeTopic, updates.str()); + } + if (m_connected && m_hasDefinitionTopic) { + deque messages; + m_messages->findAll("", "", "", false, true, true, true, true, true, 0, 0, false, &messages); + for (const auto& message : messages) { + if (message->getCreateTime() <= m_definitionsSince) { // only newer defined + continue; + } + MqttReplacers msgValues = m_replacers; // need a copy here as the contents are manipulated + msgValues.set("circuit", message->getCircuit()); + msgValues.set("name", message->getName()); + if (!m_publishByField) { + msgValues.set("topic", getTopic(message, "", "")); // TODO already present? + } + msgValues.reduce(); + ostringstream fields; + for (size_t index = 0; index < message->getFieldCount(); index++) { + const SingleDataField* field = message->getField(index); + if (!field || field->isIgnored()) { + continue; + } + const DataType* dataType = field->getDataType(); + string typeSuffix; + if (dataType->isNumeric()) { + typeSuffix = "number"; + } else { + typeSuffix = "string"; + } + //TODO valuelists => binary_sensor, date/time => device_class=date/timestamp, energy, power, current, gas, humidity, power_factor, pressure, temperature, voltage + string str = msgValues.get("type_"+typeSuffix, false, false); + if (str.empty()) { + continue; + } + MqttReplacers values = msgValues; // need a copy here as the contents are manipulated + values.set("type", str); + values.set("index", static_cast(index)); + string fieldName = message->getFieldName(index); + values.set("field", fieldName); + values.set("fieldcomment", field->getAttribute("comment")); + values.set("unit", field->getAttribute("unit")); + if (dataType->isNumeric()) { + auto dt = dynamic_cast(dataType); + values.set("min", static_cast(dt->getMinValue())); + values.set("max", static_cast(dt->getMaxValue())); + } + values.reduce(); + str = values.get("type_part_"+typeSuffix, false, false); + values.set("type_part", str); + //TODO Valuelist: => select.options + if (m_publishByField) { + values.set("topic", getTopic(message, "", fieldName)); // TODO already present? + } + values.reduce(); + if (m_hasDefinitionFieldsPayload) { + string value = values["field_payload"]; + if (!value.empty()) { + if (fields.tellp()>0) { + fields << values["field_separator"]; + } + fields << value; + // TODO str << ",\"value_template\":\"{{value_json." << fieldName << ".value}}\""; + } + continue; + } + string topic = values.get("definition_topic", false, false); + if (!topic.empty()) { + values.set("definition_topic", topic); + string payload = values.get("definition_payload", false, false); + string retainStr = values.get("definition_retain", false, false); + bool retain = !retainStr.empty() && !(retainStr=="0" || retainStr=="no" || retainStr=="false"); + publishTopic(topic, payload, retain); + } + } + if (fields.tellp()>0) { + msgValues.set("fields_payload", fields.str()); + string topic = msgValues.get("definition_topic", false, false); + if (!topic.empty()) { + msgValues.set("definition_topic", topic); + string payload = msgValues.get("definition_payload", false, false); + string retainStr = msgValues.get("definition_retain", false, false); + bool retain = !retainStr.empty() && !(retainStr == "0" || retainStr == "no" || retainStr == "false"); + publishTopic(topic, payload, retain); + } + } + } + time(&m_definitionsSince); + } time(&lastTaskRun); } if (sendSignal) { @@ -795,7 +1155,7 @@ void MqttHandler::run() { for (auto it = m_updatedMessages.begin(); it != m_updatedMessages.end(); ) { const vector* messages = m_messages->getByKey(it->first); if (messages) { - for (auto message : *messages) { + for (const auto& message : *messages) { if (message->getLastChangeTime() > 0 && message->isAvailable() && (!g_onlyChanges || message->getLastChangeTime() > lastUpdates)) { updates.str(""); @@ -871,24 +1231,15 @@ bool MqttHandler::handleTraffic(bool allowReconnect) { } string MqttHandler::getTopic(const Message* message, const string& suffix, const string& fieldName) { - ostringstream ret; - for (size_t i = 0; i < g_topicStrs.size(); i++) { - ret << g_topicStrs[i]; - if (!message) { - break; - } - if (i < g_topicFields.size()) { - if (g_topicFields[i] == "field") { - ret << fieldName; - } else { - message->dumpField(g_topicFields[i], false, OF_NONE, &ret); - } + map values; + if (message) { + values["circuit"] = message->getCircuit(); + values["name"] = message->getName(); + if (!fieldName.empty()) { + values["field"] = fieldName; } } - if (!suffix.empty()) { - ret << suffix; - } - return ret.str(); + return g_topicReplacer->get(values) + suffix; } void MqttHandler::publishMessage(const Message* message, ostringstream* updates, bool includeWithoutData) { diff --git a/src/ebusd/mqtthandler.h b/src/ebusd/mqtthandler.h index 6f35396d..f9bf7bc7 100644 --- a/src/ebusd/mqtthandler.h +++ b/src/ebusd/mqtthandler.h @@ -55,6 +55,166 @@ const struct argp_child* mqtthandler_getargs(); bool mqtthandler_register(UserInfo* userInfo, BusHandler* busHandler, MessageMap* messages, list* handlers); +/** + * Helper class for replacing a template string with real values. + */ +class MqttReplacer { + public: + /** + * Constructor. + * @param fillDefault true to fill in the default topic. + */ + explicit MqttReplacer(bool fillDefault = true); + + /** + * Create a new replacer. + * @param templateStr the template string. + * @param ensureDefault true to ensure the default topic parts are present. + * @param onlyKnown true to allow only known field names from @a knownFieldNames. + * @param noKnownDuplicates true to now allow duplicates from @a knownFieldNames. + * @return the replacer, or null if the parameters are invalid. + */ + static MqttReplacer* create(const string& templateStr, bool ensureDefault = true, bool onlyKnown = true, bool noKnownDuplicates = true); + + /** + * Normalize the string to contain only alpha numeric characters plus underscore by replacing other characters with + * an underscore. + * @param str the string to normalize. + */ + static void normalize(string& str); + + /** + * Parse the template string. + * @param templateStr the template string. + * @param onlyKnown true to allow only known field names from @a knownFieldNames. + * @param noKnownDuplicates true to now allow duplicates from @a knownFieldNames. + * @return true on success, false on malformed template string. + */ + bool parse(const string& templateStr, bool onlyKnown = true, bool noKnownDuplicates = true); + + private: + /** + * Ensure the default topic parts are present (circuit and message). + */ + void ensureDefault(); + + public: + /** + * Return whether the specified field is used. + * @param field the field name to check. + * @return true when the specified field is used + */ + bool has(const string& field) const; + + /** + * Get the replaced template string. + * @param values the named values for replacement. + * @param untilFirstEmpty true to only return the prefix before the first empty field. + * @param onlyAlphanum whether to only allow alpha numeric characters plus underscore. + * @return the replaced template string. + */ + string get(const map& values, bool untilFirstEmpty = true, bool onlyAlphanum = false) const; + + /** + * Check if the fields can be reduced to a constant value. + * @param values the named values for replacement. + * @return true if the result is final. + */ + bool isReducable(const map& values) const; + + /** + * Reduce the fields to a constant value if possible. + * @param values the named values for replacement. + * @param result the string to store the result in. + * @param onlyAlphanum whether to only allow alpha numeric characters plus underscore. + * @return true if the result is final. + */ + bool reduce(const map& values, string& result, bool onlyAlphanum = false) const; + + /** + * + * @param remain + * @param circuit + * @param name + * @return the index of the last unmatched part, or the negative index minus one for extra non-matched non-field parts. + */ + ssize_t matchTopic(const string& remain, string* circuit, string* name, string* field) const; + + private: + /** + * the list of parts the template is composed of. + * the string is either the plain string or the name of the field. + * the number is negative for plain strings, the index to @a knownFieldNames for a known field, or the size of + * @a knownFieldNames for an unknown field. + */ + vector> m_parts; +}; + + +/** + * A set of constants and @a MqttReplacer variables. + */ +class MqttReplacers { + public: + /** + * Get the value of the specified key from the constants only. + * @param key the key for which to get the value. + * @return the value string or empty. + */ + const string& operator[](const string& key) const; + + /** + * Check if the specified field is used by one of the replacers. + * @param field the name of the field to check. + * @return true if the specified field is used by one of the replacers. + */ + bool uses(const string& field) const; + + /** + * Get the variable value of the specified key. + * @param key the key for which to get the value. + * @return the value @a MqttReplacer. + */ + MqttReplacer& get(const string& key); + + /** + * Get the variable or constant value of the specified key. + * @param key the key for which to get the value. + * @param value the value string or empty. + */ + string get(const string& key, bool untilFirstEmpty, bool onlyAlphanum = false) const; + + /** + * Set the constant value of the specified key and additionally normalized with uppercase key only (if the key does + * not contain an underscore). + * @param key the key to store. + * @param value the value string. + * @param removeReplacer true to remove a replacer with the same name. + * @return true when an upper case key was stored/updates as well. + */ + bool set(const string& key, const string& value, bool removeReplacer = true); + + /** + * Set the constant value of the specified key. + * @param key the key to store. + * @param value the numeric value (converted to a string). + */ + void set(const string& key, int value); + + /** + * Reduce as many variables to constants as possible. + */ + void reduce(); + + private: + /** constant values from the integration file. */ + map m_constants; + + /** variable values from the integration file. */ + map m_replacers; +}; + + /** * The main class supporting MQTT data handling. */ @@ -68,6 +228,10 @@ class MqttHandler : public DataSink, public DataSource, public WaitThread { */ MqttHandler(UserInfo* userInfo, BusHandler* busHandler, MessageMap* messages); + private: + void parseIntegration(const string& line); + + public: /** * Destructor. */ @@ -150,6 +314,24 @@ class MqttHandler : public DataSink, public DataSource, public WaitThread { /** whether to publish a separate topic for each message field. */ bool m_publishByField; + /** the @a MqttReplacers from the integration file. */ + MqttReplacers m_replacers; + + /** whether the @a m_replacers includes the definition_topic. */ + bool m_hasDefinitionTopic; + + /** whether the @a m_replacers uses the fields_payload variable. */ + bool m_hasDefinitionFieldsPayload; + + /** the subscribed configuration restart topic, or empty. */ + string m_subscribeConfigRestartTopic; + + /** the expected payload of the subscribed configuration restart topic, or empty for any. */ + string m_subscribeConfigRestartPayload; + + /** the last system time when the message definitions were published. */ + time_t m_definitionsSince; + /** the mosquitto structure if initialized, or nullptr. */ struct mosquitto* m_mosquitto;