diff --git a/README.md b/README.md index d54bc1cd..8390a29e 100644 --- a/README.md +++ b/README.md @@ -2,7 +2,7 @@ ## Description -MQTT module for the [OpenDAQ SDK](https://github.com/openDAQ/openDAQ). The module is designed for software communication via the *MQTT 3.1.1* protocol using an external broker. It allows publishing and subscribing to openDAQ signal data over MQTT. The module consists of five key openDAQ components: the *MQTT client function block* (**MQTTClientFB**) and its nested function blocks — the *publisher* (**MQTTJSONPublisherFB**) and the *subscriber* (**MQTTSubscriberFB**) with its nested block *JSON decoder* (**MQTTJSONDecoderFB**). +MQTT module for the [OpenDAQ SDK](https://github.com/openDAQ/openDAQ). The module is designed for software communication via the *MQTT 3.1.1* protocol using an external broker. It allows publishing and subscribing to openDAQ signal data over MQTT. The module consists of five key openDAQ components: the *MQTT client device* (**OpenDAQMQTTDevice**) and its nested function blocks — the *publisher* (**MQTTJSONPublisherFB**) and the *subscriber* (**MQTTSubscriberFB**) with its nested block *JSON decoder* (**MQTTJSONDecoderFB**). ### Functional - Connecting to an MQTT broker; @@ -12,15 +12,16 @@ MQTT module for the [OpenDAQ SDK](https://github.com/openDAQ/openDAQ). The modul - A set of examples and *gtests* for verifying functionality. ### Key components -1) **MQTT client Function Block (MQTTClientFB)**: - - **Where**: *mqtt_streaming_module/src/mqtt_client_fb_impl.cpp, include/mqtt_streaming_module/...* - - **Purpose**: Represents the MQTT broker as an openDAQ function block - the connection point through which function blocks are created. +1) **MQTT client Device (OpenDAQMQTTDevice)**: + - **Where**: *mqtt_streaming_module/src/mqtt_client_device_impl.cpp, include/mqtt_streaming_module/...* + - **Purpose**: Represents the MQTT broker as an openDAQ device - the connection point through which function blocks are created. + - **Connection string**: `daq.mqtt://[:]`, for example `daq.mqtt://127.0.0.1:1883`. The device is added with `instance.addDevice(...)`. The broker host always comes from the connection string; a port given there takes precedence over the *Port* property. - **Main properties:** - - *BrokerAddress* (string) — MQTT broker address. It can be an IP address or a hostname. By default, it is set to *"127.0.0.1"*. - - *BrokerPort* (integer) — Port number for the MQTT broker connection. By default, it is set to *1883*. + - *Port* (integer) — Port number for the MQTT broker connection, used only when the connection string carries no port. By default, it is set to *1883*. - *Username* (string) — Username for MQTT broker authentication. By default, it is empty. - *Password* (string) — Password for MQTT broker authentication. By default, it is empty. - *ConnectionTimeout* (integer) — Timeout in milliseconds for the initial connection to the MQTT broker. If the connection fails, an exception is thrown. By default, it is set to *3000 ms*. + - **Connection status**: exposed through the device's connection status container under the *ConfigurationStatus* alias, e.g. `device.getConnectionStatusContainer().getStatus("ConfigurationStatus")`. 2) **MQTT publisher Function Block (MQTTJSONPublisherFB)**: - **Where**: *include/mqtt_streaming_module/mqtt_publisher_fb_impl.h, src/mqtt_publisher_fb_impl.cpp* - **Purpose**: Publishes openDAQ signal data to MQTT topics. @@ -192,15 +193,15 @@ cmake --build . ## Examples There are 3 example C++ application: - - **custom-mqtt-sub** - demonstrates how to work with the *MQTT subscriber MQTT FB* and *MQTT JSON decoder MQTT FB*. The application creates an *MQTTClientFB* and a *MQTTSubscriberFB* with nested *MQTTJSONDecoderFB* function blocks to receive JSON MQTT messages, parse them, and create openDAQ signals to send the parsed data. The application also creates *packet readers* for all FB signals and prints the samples to standard output. The *JSONConfigFile* property of the *MQTTSubscriberFB* is set to the value of path whose is provided as a command-line argument when the application starts (see the **Key components** section). Usage: + - **custom-mqtt-sub** - demonstrates how to work with the *MQTT subscriber MQTT FB* and *MQTT JSON decoder MQTT FB*. The application creates an *OpenDAQMQTTDevice* and a *MQTTSubscriberFB* with nested *MQTTJSONDecoderFB* function blocks to receive JSON MQTT messages, parse them, and create openDAQ signals to send the parsed data. The application also creates *packet readers* for all FB signals and prints the samples to standard output. The *JSONConfigFile* property of the *MQTTSubscriberFB* is set to the value of path whose is provided as a command-line argument when the application starts (see the **Key components** section). Usage: ```bash ./custom-mqtt-sub --address broker.emqx.io examples/custom-mqtt-sub/public-example0.json ``` - - **raw-mqtt-sub** - demonstrates how to work with the *MQTT subscriber MQTT FB* in a raw mode (binary data without parsing). The application creates an *MQTTClientFB* and a *MQTTSubscriberFB* to receive MQTT messages and create openDAQ signals to send the data as binary packets. The application also creates packet readers for all FB signals and prints the binary packets as strings to standard output. The *Topic* property of the *MQTTSubscriberFB* is filled from the application arguments. Usage: + - **raw-mqtt-sub** - demonstrates how to work with the *MQTT subscriber MQTT FB* in a raw mode (binary data without parsing). The application creates an *OpenDAQMQTTDevice* and a *MQTTSubscriberFB* to receive MQTT messages and create openDAQ signals to send the data as binary packets. The application also creates packet readers for all FB signals and prints the binary packets as strings to standard output. The *Topic* property of the *MQTTSubscriberFB* is filled from the application arguments. Usage: ```bash ./raw-mqtt-sub --address broker.emqx.io /mirip/UNet3AC2/sensor/data ``` - - **ref-dev-mqtt-pub** - demonstrates how to work with the *MQTTJSONPublisherFB*. The application creates an *openDAQ ref-device* with four channels, an *MQTTClientFB*, and a *MQTTJSONPublisherFB* to publish JSON MQTT messages with the channels’ data. The properties of the *MQTTJSONPublisherFB* are set according to the selected mode, which can be specified via the *--mode* option, and array size, which can be specified via the *--array* option with size. Posible *--mode* option values are: + - **ref-dev-mqtt-pub** - demonstrates how to work with the *MQTTJSONPublisherFB*. The application creates an *openDAQ ref-device* with four channels, an *OpenDAQMQTTDevice*, and a *MQTTJSONPublisherFB* to publish JSON MQTT messages with the channels’ data. The properties of the *MQTTJSONPublisherFB* are set according to the selected mode, which can be specified via the *--mode* option, and array size, which can be specified via the *--array* option with size. Posible *--mode* option values are: - 0 - One MQTT topic per signal; - 1 - One MQTT message/topic for all signals. ```bash diff --git a/examples/custom-mqtt-sub/src/custom-mqtt-sub.cpp b/examples/custom-mqtt-sub/src/custom-mqtt-sub.cpp index 9dee43a1..0db7397a 100644 --- a/examples/custom-mqtt-sub/src/custom-mqtt-sub.cpp +++ b/examples/custom-mqtt-sub/src/custom-mqtt-sub.cpp @@ -122,13 +122,11 @@ int main(int argc, char* argv[]) return appConfig.error; } - // Create OpenDAQ instance and add MQTT broker FB + // Create OpenDAQ instance and add the MQTT broker device const InstancePtr instance = InstanceBuilder().addModulePath(MODULE_PATH).build(); - const std::string clientFbName = "MQTTClientFB"; - auto clientFbConfig = instance.getAvailableFunctionBlockTypes().get(clientFbName).createDefaultConfig(); - clientFbConfig.setPropertyValue("BrokerAddress", appConfig.brokerAddress); - auto brokerFB = instance.addFunctionBlock(clientFbName, clientFbConfig); - auto availableFbs = brokerFB.getAvailableFunctionBlockTypes(); + const std::string connectionString = "daq.mqtt://" + appConfig.brokerAddress; + auto brokerDevice = instance.addDevice(connectionString); + auto availableFbs = brokerDevice.getAvailableFunctionBlockTypes(); const std::string subFbName = "MQTTSubscriberFB"; std::cout << "Try to add the " << subFbName << std::endl; @@ -137,7 +135,7 @@ int main(int argc, char* argv[]) config.setPropertyValue("JSONConfigFile", appConfig.configFilePath); // Add the subscriber function block to the broker FB - daq::FunctionBlockPtr subFb = brokerFB.addFunctionBlock(subFbName, config); + daq::FunctionBlockPtr subFb = brokerDevice.addFunctionBlock(subFbName, config); // Create packet readers for all signals auto signals = List(); diff --git a/examples/raw-mqtt-sub/src/raw-mqtt-sub.cpp b/examples/raw-mqtt-sub/src/raw-mqtt-sub.cpp index ee8398cb..ff851fac 100644 --- a/examples/raw-mqtt-sub/src/raw-mqtt-sub.cpp +++ b/examples/raw-mqtt-sub/src/raw-mqtt-sub.cpp @@ -54,13 +54,11 @@ int main(int argc, char* argv[]) return appConfig.error; } - // Create OpenDAQ instance and add MQTT broker FB + // Create OpenDAQ instance and add the MQTT broker device const InstancePtr instance = InstanceBuilder().addModulePath(MODULE_PATH).build(); - const std::string clientFbName = "MQTTClientFB"; - auto clientFbConfig = instance.getAvailableFunctionBlockTypes().get(clientFbName).createDefaultConfig(); - clientFbConfig.setPropertyValue("BrokerAddress", appConfig.brokerAddress); - auto brokerFB = instance.addFunctionBlock(clientFbName, clientFbConfig); - auto availableFbs = brokerFB.getAvailableFunctionBlockTypes(); + const std::string connectionString = "daq.mqtt://" + appConfig.brokerAddress; + auto brokerDevice = instance.addDevice(connectionString); + auto availableFbs = brokerDevice.getAvailableFunctionBlockTypes(); const std::string fbName = "MQTTSubscriberFB"; std::cout << "Try to add the " << fbName << std::endl; @@ -72,7 +70,7 @@ int main(int argc, char* argv[]) config.setPropertyValue("MessageIsString", True); // Add the subscriber function block to the broker FB - daq::FunctionBlockPtr subFb = brokerFB.addFunctionBlock(fbName, config); + daq::FunctionBlockPtr subFb = brokerDevice.addFunctionBlock(fbName, config); // Create packet readers for a signal const auto signal = subFb.getSignals()[0]; diff --git a/examples/ref-dev-mqtt-pub/src/ref-dev-mqtt-pub.cpp b/examples/ref-dev-mqtt-pub/src/ref-dev-mqtt-pub.cpp index b54579b6..fb272f30 100644 --- a/examples/ref-dev-mqtt-pub/src/ref-dev-mqtt-pub.cpp +++ b/examples/ref-dev-mqtt-pub/src/ref-dev-mqtt-pub.cpp @@ -100,12 +100,10 @@ int main(int argc, char* argv[]) channels[3].setPropertyValue("SampleRate", 100); channels[3].setPropertyValue("Frequency", 20); - // Create and configure MQTT server - const std::string clientFbName = "MQTTClientFB"; - auto clientFbConfig = instance.getAvailableFunctionBlockTypes().get(clientFbName).createDefaultConfig(); - clientFbConfig.setPropertyValue("BrokerAddress", appConfig.brokerAddress); - auto brokerFB = instance.addFunctionBlock(clientFbName, clientFbConfig); - auto availableFbs = brokerFB.getAvailableFunctionBlockTypes(); + // Connect to the MQTT broker + const std::string connectionString = "daq.mqtt://" + appConfig.brokerAddress; + auto brokerDevice = instance.addDevice(connectionString); + auto availableFbs = brokerDevice.getAvailableFunctionBlockTypes(); const std::string fbName = "MQTTJSONPublisherFB"; std::cout << "Try to add the " << fbName << std::endl; @@ -120,7 +118,7 @@ int main(int argc, char* argv[]) // Add the publisher function block to the broker device - daq::FunctionBlockPtr fb = brokerFB.addFunctionBlock(fbName, config); + daq::FunctionBlockPtr fb = brokerDevice.addFunctionBlock(fbName, config); const auto signals = refDevice.getSignals(search::Recursive(search::Any())); for (const auto& s : signals) { diff --git a/modules/mqtt_streaming_module/include/mqtt_streaming_module/constants.h b/modules/mqtt_streaming_module/include/mqtt_streaming_module/constants.h index 961a5c2c..67b5bb18 100644 --- a/modules/mqtt_streaming_module/include/mqtt_streaming_module/constants.h +++ b/modules/mqtt_streaming_module/include/mqtt_streaming_module/constants.h @@ -24,8 +24,7 @@ static constexpr const char* DEFAULT_VALUE_SIGNAL_LOCAL_ID = "MQTTValueSignal"; static constexpr const char* DEFAULT_TS_SIGNAL_LOCAL_ID = "MQTTTimestampSignal"; -static constexpr const char* PROPERTY_NAME_CLIENT_BROKER_ADDRESS = "BrokerAddress"; -static constexpr const char* PROPERTY_NAME_CLIENT_BROKER_PORT = "BrokerPort"; +static constexpr const char* PROPERTY_NAME_CLIENT_BROKER_PORT = "Port"; static constexpr const char* PROPERTY_NAME_CLIENT_USERNAME = "Username"; static constexpr const char* PROPERTY_NAME_CLIENT_PASSWORD = "Password"; static constexpr const char* PROPERTY_NAME_CLIENT_CONNECT_TIMEOUT = "ConnectionTimeout"; @@ -59,22 +58,23 @@ static constexpr const char* PUB_PREVIEW_SIGNAL_NAME = "PreviewSignal"; static constexpr const char* SUB_FB_NAME = "MQTTSubscriberFB"; static constexpr const char* PUB_FB_NAME = "MQTTJSONPublisherFB"; -static constexpr const char* CLIENT_FB_NAME = "MQTTClientFB"; static constexpr const char* JSON_DECODER_FB_NAME = "MQTTJSONDecoderFB"; -static const char* MQTT_LOCAL_CLIENT_FB_ID_PREFIX = "MQTTClientFB"; +static constexpr const char* CLIENT_DEVICE_TYPE_ID = "OpenDAQMQTTDevice"; +static constexpr const char* CLIENT_DEVICE_TYPE_NAME = "MQTT client device"; +static constexpr const char* CLIENT_DEVICE_CONN_PREFIX = "daq.mqtt"; + +static const char* MQTT_LOCAL_CLIENT_DEVICE_ID_PREFIX = "MQTTClientDevice"; static const char* MQTT_LOCAL_PUB_FB_ID_PREFIX = "MQTTJSONPublisherFB"; static const char* MQTT_LOCAL_SUB_FB_ID_PREFIX = "MQTTSubscriberFB"; static const char* MQTT_LOCAL_JSON_DECODER_FB_ID_PREFIX = "MQTTJSONDecoderFB"; -static const char* MQTT_CLIENT_FB_CON_STATUS_TYPE = "DAQ_MQTT_ConnectionStatusType"; static const char* MQTT_PUB_FB_SIG_STATUS_TYPE = "DAQ_MQTT_SignalStatusType"; static const char* MQTT_PUB_FB_PUB_STATUS_TYPE = "DAQ_MQTT_PublishingStatusType"; static const char* MQTT_PUB_FB_SET_STATUS_TYPE = "DAQ_MQTT_SettingStatusType"; -static const char* MQTT_CLIENT_FB_CON_STATUS_NAME = "ConnectionStatus"; static const char* MQTT_PUB_FB_SIG_STATUS_NAME = "SignalStatus"; static const char* MQTT_PUB_FB_PUB_STATUS_NAME = "PublishingStatus"; static const char* MQTT_PUB_FB_SET_STATUS_NAME = "SettingStatus"; diff --git a/modules/mqtt_streaming_module/include/mqtt_streaming_module/mqtt_client_fb_impl.h b/modules/mqtt_streaming_module/include/mqtt_streaming_module/mqtt_client_device_impl.h similarity index 65% rename from modules/mqtt_streaming_module/include/mqtt_streaming_module/mqtt_client_fb_impl.h rename to modules/mqtt_streaming_module/include/mqtt_streaming_module/mqtt_client_device_impl.h index 3fbc2fa2..03bd9599 100644 --- a/modules/mqtt_streaming_module/include/mqtt_streaming_module/mqtt_client_fb_impl.h +++ b/modules/mqtt_streaming_module/include/mqtt_streaming_module/mqtt_client_device_impl.h @@ -19,30 +19,32 @@ #include "mqtt_streaming_protocol/MqttSettings.h" #include #include -#include -#include -#include +#include BEGIN_NAMESPACE_OPENDAQ_MQTT_STREAMING_MODULE -class MqttClientFbImpl : public FunctionBlock +class MqttClientDeviceImpl : public Device { public: - explicit MqttClientFbImpl(const ContextPtr& ctx, - const ComponentPtr& parent, - const PropertyObjectPtr& config); + explicit MqttClientDeviceImpl(const ContextPtr& ctx, + const ComponentPtr& parent, + const StringPtr& localId, + const StringPtr& connectionString, + const std::string& brokerHost, + const PropertyObjectPtr& config); - static FunctionBlockTypePtr CreateType(); + static DeviceTypePtr CreateType(); + static PropertyObjectPtr CreateDefaultConfig(); protected: - static std::atomic localIndex; - static std::string generateLocalId(); - void removed() override; + DeviceInfoPtr onGetInfo() override; + DictPtr onGetAvailableFunctionBlockTypes() override; FunctionBlockPtr onAddFunctionBlock(const StringPtr& typeId, const PropertyObjectPtr& config) override; + void onRemoveFunctionBlock(const FunctionBlockPtr& functionBlock) override; void initNestedFbTypes(); void initMqttSubscriber(); @@ -51,10 +53,14 @@ class MqttClientFbImpl : public FunctionBlock void readProperties(const PropertyObjectPtr& config); bool waitForConnection(const int timeoutMs); + /// Pushes a ConnectionStatusType value into the device's connection status container. + /// Surfaces to clients under the "ConfigurationStatus" alias. + void setConnectionStatus(const std::string& value, const std::string& message = ""); + DictObjectPtr nestedFbTypes; + StringPtr connectionString; int connectTimeout; - StatusAdaptor connectionStatus; std::shared_ptr subscriber; Mqtt::Utils::Settings::MqttConnectionSettings connectionSettings; diff --git a/modules/mqtt_streaming_module/include/mqtt_streaming_module/mqtt_streaming_module_impl.h b/modules/mqtt_streaming_module/include/mqtt_streaming_module/mqtt_streaming_module_impl.h index 70aecf22..72a1825a 100644 --- a/modules/mqtt_streaming_module/include/mqtt_streaming_module/mqtt_streaming_module_impl.h +++ b/modules/mqtt_streaming_module/include/mqtt_streaming_module/mqtt_streaming_module_impl.h @@ -18,6 +18,7 @@ #include #include #include +#include BEGIN_NAMESPACE_OPENDAQ_MQTT_STREAMING_MODULE @@ -27,14 +28,29 @@ class MqttStreamingModule final : public Module public: MqttStreamingModule(ContextPtr context); - DictPtr onGetAvailableFunctionBlockTypes() override; - FunctionBlockPtr onCreateFunctionBlock(const StringPtr& id, - const ComponentPtr& parent, - const StringPtr& localId, - const PropertyObjectPtr& config) override; + DictPtr onGetAvailableDeviceTypes() override; + DevicePtr onCreateDevice(const StringPtr& connectionString, + const ComponentPtr& parent, + const PropertyObjectPtr& config) override; + + /// Host and port taken from a `daq.mqtt://host[:port]` connection string. + struct BrokerAddress + { + std::string host; + uint16_t port{0}; // 0 when the connection string carries no port + }; + + /// Throws InvalidParameterException when the string is not a valid `daq.mqtt://` address. + DAQ_MQTT_STREAM_MODULE_API static BrokerAddress parseConnectionString(const StringPtr& connectionString); + DAQ_MQTT_STREAM_MODULE_API static StringPtr formatConnectionString(const MqttStreamingModule::BrokerAddress& conParam); private: - static FunctionBlockTypePtr createFbType(); + static DeviceTypePtr createDeviceType(); + + static PropertyObjectPtr populateDefaultConfig(const PropertyObjectPtr& config); + + std::mutex sync; + size_t deviceIndex{0}; }; END_NAMESPACE_OPENDAQ_MQTT_STREAMING_MODULE diff --git a/modules/mqtt_streaming_module/include/mqtt_streaming_module/status_adaptor.h b/modules/mqtt_streaming_module/include/mqtt_streaming_module/status_adaptor.h deleted file mode 100644 index 2c0a4fed..00000000 --- a/modules/mqtt_streaming_module/include/mqtt_streaming_module/status_adaptor.h +++ /dev/null @@ -1,72 +0,0 @@ -/* - * Copyright 2022-2025 openDAQ d.o.o. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -#pragma once -#include -#include -#include -#include - -BEGIN_NAMESPACE_OPENDAQ_MQTT_STREAMING_MODULE - -class StatusAdaptor -{ -public: - StatusAdaptor(const std::string typeName, - const std::string statusName, - ComponentStatusContainerPtr statusContainer, - std::string initState, - TypeManagerPtr typeManager) - : typeName(typeName), - statusName(statusName), - statusContainer(statusContainer), - typeManager(typeManager) - { - currentStatus = Enumeration(typeName, initState, typeManager); - currentMessage = ""; - statusContainer.template asPtr(true).addStatus(statusName, currentStatus); - } - - bool setStatus(const std::string& status, const std::string& message = "") - { - std::scoped_lock lock(statusMutex); - const auto newStatus = Enumeration(typeName, String(status), typeManager); - bool changed = (newStatus != currentStatus || message != currentMessage); - if (changed) - { - currentStatus = newStatus; - currentMessage = message; - statusContainer.template asPtr(true).setStatusWithMessage(statusName, currentStatus, message); - } - return changed; - } - std::string getStatus() - { - std::scoped_lock lock(statusMutex); - return currentStatus.getValue().toStdString(); - } - -private: - const std::string typeName; - const std::string statusName; - std::string currentMessage; - ComponentStatusContainerPtr statusContainer; - EnumerationPtr currentStatus; - TypeManagerPtr typeManager; - std::mutex statusMutex; -}; - -END_NAMESPACE_OPENDAQ_MQTT_STREAMING_MODULE diff --git a/modules/mqtt_streaming_module/src/CMakeLists.txt b/modules/mqtt_streaming_module/src/CMakeLists.txt index e61ecdba..2dfb25af 100644 --- a/modules/mqtt_streaming_module/src/CMakeLists.txt +++ b/modules/mqtt_streaming_module/src/CMakeLists.txt @@ -4,7 +4,7 @@ set(MODULE_HEADERS_DIR ../include/${TARGET_FOLDER_NAME}) set(SRC_Include common.h module_dll.h mqtt_streaming_module_impl.h - mqtt_client_fb_impl.h + mqtt_client_device_impl.h mqtt_json_decoder_fb_impl.h mqtt_subscriber_fb_impl.h mqtt_publisher_fb_impl.h @@ -21,13 +21,12 @@ set(SRC_Include common.h types.h status_helper.h property_helper.h - status_adaptor.h status_container.h ) set(SRC_Srcs module_dll.cpp mqtt_streaming_module_impl.cpp - mqtt_client_fb_impl.cpp + mqtt_client_device_impl.cpp mqtt_json_decoder_fb_impl.cpp mqtt_subscriber_fb_impl.cpp mqtt_publisher_fb_impl.cpp @@ -48,7 +47,6 @@ source_group("common" FILES ${MODULE_HEADERS_DIR}/common.h ${MODULE_HEADERS_DIR}/types.h ${MODULE_HEADERS_DIR}/status_helper.h ${MODULE_HEADERS_DIR}/property_helper.h - ${MODULE_HEADERS_DIR}/status_adaptor.h ${MODULE_HEADERS_DIR}/status_container.h helper.cpp ) @@ -65,8 +63,10 @@ source_group("functionalBlock" FILES ${MODULE_HEADERS_DIR}/mqtt_json_decoder_fb_ mqtt_subscriber_fb_impl.cpp ${MODULE_HEADERS_DIR}/mqtt_publisher_fb_impl.h mqtt_publisher_fb_impl.cpp - ${MODULE_HEADERS_DIR}/mqtt_client_fb_impl.h - mqtt_client_fb_impl.cpp +) + +source_group("device" FILES ${MODULE_HEADERS_DIR}/mqtt_client_device_impl.h + mqtt_client_device_impl.cpp ) source_group("handlers" FILES ${MODULE_HEADERS_DIR}/handler_base.h diff --git a/modules/mqtt_streaming_module/src/mqtt_client_fb_impl.cpp b/modules/mqtt_streaming_module/src/mqtt_client_device_impl.cpp similarity index 71% rename from modules/mqtt_streaming_module/src/mqtt_client_fb_impl.cpp rename to modules/mqtt_streaming_module/src/mqtt_client_device_impl.cpp index cba688f5..f9426b93 100644 --- a/modules/mqtt_streaming_module/src/mqtt_client_fb_impl.cpp +++ b/modules/mqtt_streaming_module/src/mqtt_client_device_impl.cpp @@ -1,7 +1,10 @@ #include "mqtt_streaming_module/constants.h" #include "mqtt_streaming_module/mqtt_subscriber_fb_impl.h" #include "mqtt_streaming_module/mqtt_publisher_fb_impl.h" -#include +#include +#include +#include +#include #include #include @@ -9,21 +12,25 @@ BEGIN_NAMESPACE_OPENDAQ_MQTT_STREAMING_MODULE constexpr int MQTT_CLIENT_SYNC_DISCONNECT_TOUT = 3000; -std::atomic MqttClientFbImpl::localIndex = 0; - -MqttClientFbImpl::MqttClientFbImpl(const ContextPtr& ctx, const ComponentPtr& parent, const PropertyObjectPtr& config) - : FunctionBlock(CreateType(), ctx, parent, generateLocalId()), +MqttClientDeviceImpl::MqttClientDeviceImpl(const ContextPtr& ctx, + const ComponentPtr& parent, + const StringPtr& localId, + const StringPtr& connectionString, + const std::string& brokerHost, + const PropertyObjectPtr& config) + : Device(ctx, parent, localId, nullptr, CLIENT_DEVICE_TYPE_NAME), + connectionString(connectionString), connectTimeout(0), - connectionStatus("ConnectionStatusType", - MQTT_CLIENT_FB_CON_STATUS_NAME, - statusContainer, - "Reconnecting", - context.getTypeManager()), subscriber(std::make_shared()) { initComponentStatus(); + connectionStatusContainer.addConfigurationConnectionStatus( + connectionString, Enumeration("ConnectionStatusType", "Reconnecting", context.getTypeManager())); initConnectionStatus(); initProperties(config); + + connectionSettings.mqttUrl = brokerHost; + initNestedFbTypes(); initMqttSubscriber(); @@ -36,9 +43,9 @@ MqttClientFbImpl::MqttClientFbImpl(const ContextPtr& ctx, const ComponentPtr& pa LOG_I("MQTT: Connection established"); } -void MqttClientFbImpl::removed() +void MqttClientDeviceImpl::removed() { - FunctionBlock::removed(); + Device::removed(); LOG_I("MQTT: disconnecting from the MQTT broker...", connectionSettings.mqttUrl + ":" + std::to_string(connectionSettings.port)); bool disRes = subscriber->syncDisconnect(MQTT_CLIENT_SYNC_DISCONNECT_TOUT); if (!disRes) @@ -51,7 +58,14 @@ void MqttClientFbImpl::removed() } } -void MqttClientFbImpl::initNestedFbTypes() +DeviceInfoPtr MqttClientDeviceImpl::onGetInfo() +{ + auto info = DeviceInfo(connectionString, CLIENT_DEVICE_TYPE_NAME); + info.setDeviceType(CreateType()); + return info; +} + +void MqttClientDeviceImpl::initNestedFbTypes() { nestedFbTypes = Dict(); // Add a MQTT subscriber function block type @@ -67,7 +81,7 @@ void MqttClientFbImpl::initNestedFbTypes() } } -void MqttClientFbImpl::initMqttSubscriber() +void MqttClientDeviceImpl::initMqttSubscriber() { const auto serverUrl = connectionSettings.mqttUrl + ((connectionSettings.port > 0) ? ":" + std::to_string(connectionSettings.port) : ""); subscriber->setServerURL(serverUrl); @@ -84,7 +98,7 @@ void MqttClientFbImpl::initMqttSubscriber() bool expected = false; if (connectedDone.compare_exchange_strong(expected, true)) { - connectionStatus.setStatus("Connected"); + setConnectionStatus("Connected"); connectedPromise->set_value(true); std::scoped_lock lock(componentStatusSync); setComponentStatus(ComponentStatus::Ok); @@ -95,18 +109,18 @@ void MqttClientFbImpl::initMqttSubscriber() subscriber->connect(); } -void MqttClientFbImpl::initConnectionStatus() +void MqttClientDeviceImpl::initConnectionStatus() { subscriber->setConnectionLostCb( [this](std::string msg) { - connectionStatus.setStatus("Reconnecting", msg); + setConnectionStatus("Reconnecting", msg); std::scoped_lock lock(componentStatusSync); setComponentStatusWithMessage(ComponentStatus::Error, "Connection lost"); }); } -void MqttClientFbImpl::initProperties(const PropertyObjectPtr& config) +void MqttClientDeviceImpl::initProperties(const PropertyObjectPtr& config) { for (const auto& prop : config.getAllProperties()) { @@ -145,10 +159,9 @@ void MqttClientFbImpl::initProperties(const PropertyObjectPtr& config) readProperties(config); } -void MqttClientFbImpl::readProperties(const PropertyObjectPtr& config) +void MqttClientDeviceImpl::readProperties(const PropertyObjectPtr& config) { - connectionSettings.mqttUrl = config.getPropertyValue(PROPERTY_NAME_CLIENT_BROKER_ADDRESS).asPtr().toStdString(); - connectionSettings.port = config.getPropertyValue(PROPERTY_NAME_CLIENT_BROKER_PORT); + connectionSettings.port = config.getPropertyValue(PROPERTY_NAME_CLIENT_BROKER_PORT).asPtr(); connectionSettings.username = config.getPropertyValue(PROPERTY_NAME_CLIENT_USERNAME).asPtr().toStdString(); connectionSettings.password = config.getPropertyValue(PROPERTY_NAME_CLIENT_PASSWORD).asPtr().toStdString(); connectionSettings.clientId = globalId.toStdString(); @@ -156,26 +169,32 @@ void MqttClientFbImpl::readProperties(const PropertyObjectPtr& config) connectTimeout = config.getPropertyValue(PROPERTY_NAME_CLIENT_CONNECT_TIMEOUT); } -bool MqttClientFbImpl::waitForConnection(const int timeoutMs) +bool MqttClientDeviceImpl::waitForConnection(const int timeoutMs) { bool res = (connectedFuture.wait_for(std::chrono::milliseconds(timeoutMs)) == std::future_status::ready && connectedFuture.get() == true); subscriber->setConnectedCb( [this] { - connectionStatus.setStatus("Connected"); + setConnectionStatus("Connected"); std::scoped_lock lock(componentStatusSync); setComponentStatus(ComponentStatus::Ok); }); return res; } -DictPtr MqttClientFbImpl::onGetAvailableFunctionBlockTypes() +void MqttClientDeviceImpl::setConnectionStatus(const std::string& value, const std::string& message) +{ + connectionStatusContainer.updateConnectionStatusWithMessage( + connectionString, Enumeration("ConnectionStatusType", value, context.getTypeManager()), nullptr, String(message)); +} + +DictPtr MqttClientDeviceImpl::onGetAvailableFunctionBlockTypes() { return nestedFbTypes; } -FunctionBlockPtr MqttClientFbImpl::onAddFunctionBlock(const StringPtr& typeId, const PropertyObjectPtr& config) +FunctionBlockPtr MqttClientDeviceImpl::onAddFunctionBlock(const StringPtr& typeId, const PropertyObjectPtr& config) { FunctionBlockPtr nestedFunctionBlock; { @@ -209,21 +228,20 @@ FunctionBlockPtr MqttClientFbImpl::onAddFunctionBlock(const StringPtr& typeId, c return nestedFunctionBlock; } -std::string MqttClientFbImpl::generateLocalId() +void MqttClientDeviceImpl::onRemoveFunctionBlock(const FunctionBlockPtr& functionBlock) { - return std::string(MQTT_LOCAL_CLIENT_FB_ID_PREFIX + std::to_string(localIndex++)); + auto lock = getRecursiveConfigLock2(); + + if (!functionBlocks.hasItem(functionBlock.getLocalId())) + DAQ_THROW_EXCEPTION(NotFoundException, "Function block not found: " + functionBlock.getLocalId().toStdString()); + + functionBlocks.removeItem(functionBlock); + setComponentStatus(ComponentStatus::Ok); } -FunctionBlockTypePtr MqttClientFbImpl::CreateType() +PropertyObjectPtr MqttClientDeviceImpl::CreateDefaultConfig() { auto config = PropertyObject(); - { - auto builder = - StringPropertyBuilder(PROPERTY_NAME_CLIENT_BROKER_ADDRESS, DEFAULT_BROKER_ADDRESS) - .setDescription(fmt::format("MQTT broker address. It can be an IP address or a hostname. By default it is set to \"{}\".", - DEFAULT_BROKER_ADDRESS)); - config.addProperty(builder.build()); - } { auto builder = StringPropertyBuilder(PROPERTY_NAME_CLIENT_USERNAME, DEFAULT_USERNAME) @@ -241,7 +259,9 @@ FunctionBlockTypePtr MqttClientFbImpl::CreateType() IntPropertyBuilder(PROPERTY_NAME_CLIENT_BROKER_PORT, DEFAULT_PORT) .setMinValue(1) .setMaxValue(65535) - .setDescription(fmt::format("Port number for MQTT broker connection. By default it is set to {}.", DEFAULT_PORT)); + .setDescription(fmt::format("Port number for MQTT broker connection. Used only when the connection string " + "carries no port. By default it is set to {}.", + DEFAULT_PORT)); config.addProperty(builder.build()); } { @@ -253,11 +273,15 @@ FunctionBlockTypePtr MqttClientFbImpl::CreateType() DEFAULT_INIT_TIMEOUT)); config.addProperty(builder.build()); } - const auto fbType = FunctionBlockType(CLIENT_FB_NAME, - CLIENT_FB_NAME, - "The MQTT function block allows connecting to MQTT broker. It may contain nested " - "publisher/subscriber FBs.", - config); - return fbType; + return config; +} + +DeviceTypePtr MqttClientDeviceImpl::CreateType() +{ + return DeviceType(CLIENT_DEVICE_TYPE_ID, + CLIENT_DEVICE_TYPE_NAME, + "The MQTT device connects to an MQTT broker. It may contain nested publisher/subscriber FBs.", + CLIENT_DEVICE_CONN_PREFIX, + CreateDefaultConfig()); } END_NAMESPACE_OPENDAQ_MQTT_STREAMING_MODULE diff --git a/modules/mqtt_streaming_module/src/mqtt_streaming_module_impl.cpp b/modules/mqtt_streaming_module/src/mqtt_streaming_module_impl.cpp index 5e68831a..81916b79 100644 --- a/modules/mqtt_streaming_module/src/mqtt_streaming_module_impl.cpp +++ b/modules/mqtt_streaming_module/src/mqtt_streaming_module_impl.cpp @@ -4,17 +4,22 @@ #include #include #include -#include +#include #include #include #include +#include #include -#include #include #include +#include + BEGIN_NAMESPACE_OPENDAQ_MQTT_STREAMING_MODULE +static const std::regex RegexIpv6Hostname(R"(^(.+://)?(\[[a-fA-F0-9:]+(?:\%[a-zA-Z0-9_\.-~]+)?\])(?::(\d+))?(/.*)?$)"); +static const std::regex RegexIpv4Hostname(R"(^(.+://)?([^:/\s]+)(?::(\d+))?(/.*)?$)"); + MqttStreamingModule::MqttStreamingModule(ContextPtr context) : Module(MODULE_NAME, daq::VersionInfo(MQTT_STREAM_MODULE_MAJOR_VERSION, @@ -26,34 +31,123 @@ MqttStreamingModule::MqttStreamingModule(ContextPtr context) loggerComponent = this->context.getLogger().getOrAddComponent(SHORT_MODULE_NAME); } -DictPtr MqttStreamingModule::onGetAvailableFunctionBlockTypes() +DictPtr MqttStreamingModule::onGetAvailableDeviceTypes() { - auto result = Dict(); + auto result = Dict(); - auto fbType = createFbType(); - result.set(fbType.getId(), fbType); + auto deviceType = createDeviceType(); + result.set(deviceType.getId(), deviceType); return result; } -FunctionBlockPtr -MqttStreamingModule::onCreateFunctionBlock(const StringPtr& /*id*/, - const ComponentPtr& parent, - const StringPtr& /*localId*/, - const PropertyObjectPtr& config) +DevicePtr MqttStreamingModule::onCreateDevice(const StringPtr& connectionString, + const ComponentPtr& parent, + const PropertyObjectPtr& config) { if (!context.assigned()) DAQ_THROW_EXCEPTION(InvalidParameterException, "Context is not available."); + if (!connectionString.assigned()) + DAQ_THROW_EXCEPTION(ArgumentNullException, "Connection string is not assigned."); + + PropertyObjectPtr deviceConfig = populateDefaultConfig(config); + auto conParam = parseConnectionString(connectionString); + + // A port in the connection string wins over the Port property; the property is the fallback. + if (conParam.port == 0) + { + conParam.port = static_cast(deviceConfig.getPropertyValue(PROPERTY_NAME_CLIENT_BROKER_PORT).asPtr()); + } + else + { + deviceConfig.setPropertyValue(PROPERTY_NAME_CLIENT_BROKER_PORT, conParam.port); + } + + const auto formedConnectionString = formatConnectionString(conParam); + + StringPtr localId; + { + std::scoped_lock lock(sync); + localId = String(fmt::format("{}{}", MQTT_LOCAL_CLIENT_DEVICE_ID_PREFIX, deviceIndex++)); + } - FunctionBlockPtr fb = createWithImplementation(context, parent, config); + DevicePtr device = createWithImplementation( + context, parent, localId, formedConnectionString, conParam.host, deviceConfig); - LOG_I("MQTT function block (GlobalId: {}) created", fb.getGlobalId()); + const auto deviceType = createDeviceType(); + ServerCapabilityConfigPtr connectionInfo = device.getInfo().getConfigurationConnectionInfo(); + connectionInfo.setProtocolId(deviceType.getId()); + connectionInfo.setProtocolName(deviceType.getId()); + connectionInfo.setProtocolType(ProtocolType::Unknown); + connectionInfo.setConnectionType("TCP/IP"); + connectionInfo.addAddress(conParam.host); + connectionInfo.setPort(conParam.port); + connectionInfo.setPrefix(CLIENT_DEVICE_CONN_PREFIX); + connectionInfo.setConnectionString(formedConnectionString); - return fb; + LOG_I("MQTT device (GlobalId: {}) created", device.getGlobalId()); + + return device; } -FunctionBlockTypePtr MqttStreamingModule::createFbType() +MqttStreamingModule::BrokerAddress MqttStreamingModule::parseConnectionString(const StringPtr& connectionString) { - return MqttClientFbImpl::CreateType(); + const std::string url = connectionString.toStdString(); + MqttStreamingModule::BrokerAddress conParam; + std::smatch match; + bool parsed = std::regex_search(url, match, RegexIpv6Hostname); + if (!parsed) + parsed = std::regex_search(url, match, RegexIpv4Hostname); + + if (!parsed) + DAQ_THROW_EXCEPTION(InvalidParameterException, "Could not parse connection string: {}", connectionString); + + const std::string prefix = match[1].matched ? match[1].str() : ""; + const auto expectedPrefix = std::string(CLIENT_DEVICE_CONN_PREFIX) + "://"; + if (prefix != expectedPrefix) + DAQ_THROW_EXCEPTION(InvalidParameterException, + "Connection string \"{}\" does not start with \"{}\"", + connectionString, + expectedPrefix); + + + conParam.host = match[2].str(); + if (conParam.host.empty()) + DAQ_THROW_EXCEPTION(InvalidParameterException, "Connection string \"{}\" carries no broker host", connectionString); + + if (match[3].matched) + { + const auto portValue = std::stoi(match[3].str()); + if (portValue < 1 || portValue > 65535) + DAQ_THROW_EXCEPTION(InvalidParameterException, "Port {} in connection string is out of range", portValue); + conParam.port = static_cast(portValue); + } + + return conParam; +} + +StringPtr MqttStreamingModule::formatConnectionString(const MqttStreamingModule::BrokerAddress& conParam) +{ + return String(fmt::format("{}://{}:{}", CLIENT_DEVICE_CONN_PREFIX, conParam.host, conParam.port)); +} + +DeviceTypePtr MqttStreamingModule::createDeviceType() +{ + return MqttClientDeviceImpl::CreateType(); +} + +PropertyObjectPtr MqttStreamingModule::populateDefaultConfig(const PropertyObjectPtr& config) +{ + const auto defConfig = MqttClientDeviceImpl::CreateDefaultConfig(); + if (!config.assigned()) + return defConfig; + for (const auto& prop : defConfig.getAllProperties()) + { + const auto name = prop.getName(); + if (config.hasProperty(name)) + defConfig.setPropertyValue(name, config.getPropertyValue(name)); + } + + return defConfig; } END_NAMESPACE_OPENDAQ_MQTT_STREAMING_MODULE diff --git a/modules/mqtt_streaming_module/src/mqtt_subscriber_fb_impl.cpp b/modules/mqtt_streaming_module/src/mqtt_subscriber_fb_impl.cpp index 5564769f..4a73766e 100644 --- a/modules/mqtt_streaming_module/src/mqtt_subscriber_fb_impl.cpp +++ b/modules/mqtt_streaming_module/src/mqtt_subscriber_fb_impl.cpp @@ -70,6 +70,10 @@ void MqttSubscriberFbImpl::removed() { stopProcessingThread(); unsubscribeFromTopic(); + { + auto lockProcessing = std::scoped_lock(processingMutex); + nestedFunctionBlocks.clear(); + } FunctionBlock::removed(); } diff --git a/modules/mqtt_streaming_module/tests/CMakeLists.txt b/modules/mqtt_streaming_module/tests/CMakeLists.txt index eadffd10..a2026f9a 100644 --- a/modules/mqtt_streaming_module/tests/CMakeLists.txt +++ b/modules/mqtt_streaming_module/tests/CMakeLists.txt @@ -2,7 +2,7 @@ set(MODULE_NAME mqtt_stream_module) set(TEST_APP test_${MODULE_NAME}) set(TEST_SOURCES test_mqtt_streaming_module.cpp - test_mqtt_client_fb.cpp + test_mqtt_client_device.cpp test_mqtt_subscriber_fb.cpp test_mqtt_json_decoder_fb.cpp test_mqtt_publisher_fb.cpp diff --git a/modules/mqtt_streaming_module/tests/test_daq_test_helper.h b/modules/mqtt_streaming_module/tests/test_daq_test_helper.h index a3b1eec4..17498763 100644 --- a/modules/mqtt_streaming_module/tests/test_daq_test_helper.h +++ b/modules/mqtt_streaming_module/tests/test_daq_test_helper.h @@ -2,6 +2,7 @@ #include #include #include +#include using namespace daq::modules::mqtt_streaming_module; @@ -9,13 +10,13 @@ class DaqTestHelper { public: daq::InstancePtr daqInstance; - daq::FunctionBlockPtr clientMqttFb; + daq::DevicePtr mqttDevice; daq::FunctionBlockPtr subMqttFb; void StartUp(std::string url = DEFAULT_BROKER_ADDRESS, uint16_t port = DEFAULT_PORT) { DaqInstanceInit(); - DaqMqttFbInit(url, port); + DaqMqttDeviceInit(url, port); } daq::InstancePtr DaqInstanceInit() @@ -25,32 +26,27 @@ class DaqTestHelper return daqInstance; } - daq::FunctionBlockPtr DaqAddClientMqttFb(std::string url = DEFAULT_BROKER_ADDRESS, uint16_t port = DEFAULT_PORT) + static std::string MqttConnectionString(std::string url = DEFAULT_BROKER_ADDRESS, uint16_t port = DEFAULT_PORT) { - auto config = DaqMqttFbConfig(url, port); - clientMqttFb = daqInstance.addFunctionBlock(CLIENT_FB_NAME, config); - return clientMqttFb; + return std::string(CLIENT_DEVICE_CONN_PREFIX) + "://" + url + ":" + std::to_string(port); } - daq::FunctionBlockPtr DaqMqttFbInit(std::string url = DEFAULT_BROKER_ADDRESS, uint16_t port = DEFAULT_PORT) + daq::DevicePtr DaqAddMqttDevice(std::string url = DEFAULT_BROKER_ADDRESS, uint16_t port = DEFAULT_PORT) { - if (!clientMqttFb.assigned()) - { - auto config = DaqMqttFbConfig(url, port); - clientMqttFb = daqInstance.addFunctionBlock(CLIENT_FB_NAME, config); - } - return clientMqttFb; + mqttDevice = daqInstance.addDevice(MqttConnectionString(url, port), DaqMqttDeviceConfig()); + return mqttDevice; } - daq::PropertyObjectPtr DaqMqttFbConfig(std::string url = DEFAULT_BROKER_ADDRESS, uint16_t port = DEFAULT_PORT) + daq::DevicePtr DaqMqttDeviceInit(std::string url = DEFAULT_BROKER_ADDRESS, uint16_t port = DEFAULT_PORT) { - daq::ModulePtr module; - createModule(&module, daq::NullContext()); + if (!mqttDevice.assigned()) + DaqAddMqttDevice(url, port); + return mqttDevice; + } - auto config = module.getAvailableFunctionBlockTypes().get(daq::modules::mqtt_streaming_module::CLIENT_FB_NAME).createDefaultConfig(); - config.setPropertyValue(PROPERTY_NAME_CLIENT_BROKER_ADDRESS, url); - config.setPropertyValue(PROPERTY_NAME_CLIENT_BROKER_PORT, port); - return config; + static daq::PropertyObjectPtr DaqMqttDeviceConfig() + { + return CreateModule().getAvailableDeviceTypes().get(CLIENT_DEVICE_TYPE_ID).createDefaultConfig(); } static daq::ModulePtr CreateModule() @@ -62,9 +58,9 @@ class DaqTestHelper daq::FunctionBlockPtr AddSubFb(std::string topic = "") { - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_SUB_TOPIC, daq::String(topic)); - subMqttFb = clientMqttFb.addFunctionBlock(SUB_FB_NAME, config); + subMqttFb = mqttDevice.addFunctionBlock(SUB_FB_NAME, config); return subMqttFb; } }; diff --git a/modules/mqtt_streaming_module/tests/test_mqtt_client_device.cpp b/modules/mqtt_streaming_module/tests/test_mqtt_client_device.cpp new file mode 100644 index 00000000..39b3865a --- /dev/null +++ b/modules/mqtt_streaming_module/tests/test_mqtt_client_device.cpp @@ -0,0 +1,355 @@ +#include "test_daq_test_helper.h" +#include +#include +#include +#include +#include +#include + +using namespace daq; +using namespace daq::modules::mqtt_streaming_module; + +namespace daq::modules::mqtt_streaming_module +{ +class MqttDeviceTest : public testing::Test, public DaqTestHelper +{ +}; +} // namespace daq::modules::mqtt_streaming_module + +// A port nothing is expected to listen on, used for negative connection tests. +static constexpr uint16_t UNREACHABLE_PORT = 1884; + +TEST_F(MqttDeviceTest, DefaultMqttDeviceConfig) +{ + const auto module = CreateModule(); + + DictPtr types; + ASSERT_NO_THROW(types = module.getAvailableDeviceTypes()); + ASSERT_EQ(types.getCount(), 1u); + + ASSERT_TRUE(types.hasKey(CLIENT_DEVICE_TYPE_ID)); + const auto deviceType = types.get(CLIENT_DEVICE_TYPE_ID); + ASSERT_EQ(deviceType.getId(), CLIENT_DEVICE_TYPE_ID); + ASSERT_EQ(deviceType.getName(), CLIENT_DEVICE_TYPE_NAME); + ASSERT_EQ(deviceType.getConnectionStringPrefix(), CLIENT_DEVICE_CONN_PREFIX); + + auto defaultConfig = deviceType.createDefaultConfig(); + ASSERT_TRUE(defaultConfig.assigned()); + + // BrokerAddress is gone: the host now comes from the connection string. + ASSERT_EQ(defaultConfig.getAllProperties().getCount(), 4u); + + ASSERT_FALSE(defaultConfig.hasProperty("BrokerAddress")); + ASSERT_TRUE(defaultConfig.hasProperty(PROPERTY_NAME_CLIENT_BROKER_PORT)); + ASSERT_TRUE(defaultConfig.hasProperty(PROPERTY_NAME_CLIENT_USERNAME)); + ASSERT_TRUE(defaultConfig.hasProperty(PROPERTY_NAME_CLIENT_PASSWORD)); + ASSERT_TRUE(defaultConfig.hasProperty(PROPERTY_NAME_CLIENT_CONNECT_TIMEOUT)); + + ASSERT_EQ(defaultConfig.getProperty(PROPERTY_NAME_CLIENT_BROKER_PORT).getValueType(), CoreType::ctInt); + ASSERT_EQ(defaultConfig.getProperty(PROPERTY_NAME_CLIENT_USERNAME).getValueType(), CoreType::ctString); + ASSERT_EQ(defaultConfig.getProperty(PROPERTY_NAME_CLIENT_PASSWORD).getValueType(), CoreType::ctString); + ASSERT_EQ(defaultConfig.getProperty(PROPERTY_NAME_CLIENT_CONNECT_TIMEOUT).getValueType(), CoreType::ctInt); + + ASSERT_EQ(defaultConfig.getPropertyValue(PROPERTY_NAME_CLIENT_BROKER_PORT), DEFAULT_PORT); + ASSERT_EQ(defaultConfig.getPropertyValue(PROPERTY_NAME_CLIENT_USERNAME), DEFAULT_USERNAME); + ASSERT_EQ(defaultConfig.getPropertyValue(PROPERTY_NAME_CLIENT_PASSWORD), DEFAULT_PASSWORD); + ASSERT_EQ(defaultConfig.getPropertyValue(PROPERTY_NAME_CLIENT_CONNECT_TIMEOUT), DEFAULT_INIT_TIMEOUT); +} + +TEST_F(MqttDeviceTest, ModuleExposesNoFunctionBlockTypes) +{ + const auto module = CreateModule(); + + // The client is a device now; the module no longer offers any top-level function block. + DictPtr fbTypes; + ASSERT_NO_THROW(fbTypes = module.getAvailableFunctionBlockTypes()); + ASSERT_EQ(fbTypes.getCount(), 0u); +} + +TEST_F(MqttDeviceTest, MissingPasswordProperty) +{ + const auto instance = Instance(); + DevicePtr device; + ASSERT_NO_THROW(device = instance.addDevice(MqttConnectionString())); + ASSERT_EQ(device.getStatusContainer().getStatus("ComponentStatus"), + Enumeration("ComponentStatusType", "Ok", instance.getContext().getTypeManager())); + + ASSERT_FALSE(device.hasProperty(PROPERTY_NAME_CLIENT_PASSWORD)); +} + +TEST_F(MqttDeviceTest, CreatingMqttDeviceWithDefaultConfig) +{ + const auto instance = Instance(); + DevicePtr device; + ASSERT_NO_THROW(device = instance.addDevice(MqttConnectionString())); + ASSERT_EQ(device.getStatusContainer().getStatus("ComponentStatus"), + Enumeration("ComponentStatusType", "Ok", instance.getContext().getTypeManager())); + + // The client now lives under the instance's devices, not its function blocks. + ASSERT_EQ(instance.getFunctionBlocks().getCount(), 0u); + + auto devices = instance.getDevices(); + bool contain = false; + DevicePtr deviceFromList; + for (const auto& dev : devices) + { + contain = (dev.getLocalId().toStdString().find(MQTT_LOCAL_CLIENT_DEVICE_ID_PREFIX) != std::string::npos); + if (contain) + { + deviceFromList = dev; + break; + } + } + ASSERT_TRUE(contain); + ASSERT_TRUE(deviceFromList.assigned()); + ASSERT_TRUE(deviceFromList == device); +} + +TEST_F(MqttDeviceTest, CreatingMqttDeviceWithCustomConfig) +{ + const auto instance = Instance(); + DevicePtr device; + auto config = DaqMqttDeviceConfig(); + ASSERT_NO_THROW(device = instance.addDevice(MqttConnectionString(), config)); + ASSERT_EQ(device.getStatusContainer().getStatus("ComponentStatus"), + Enumeration("ComponentStatusType", "Ok", instance.getContext().getTypeManager())); +} + +TEST_F(MqttDeviceTest, CreatingMqttDeviceWithEmptyConfig) +{ + const auto instance = Instance(); + DevicePtr device; + auto config = PropertyObject(); + ASSERT_NO_THROW(device = instance.addDevice(MqttConnectionString(), config)); + ASSERT_EQ(device.getStatusContainer().getStatus("ComponentStatus"), + Enumeration("ComponentStatusType", "Ok", instance.getContext().getTypeManager())); +} + +TEST_F(MqttDeviceTest, CreatingMqttDeviceWithPartialConfig) +{ + const auto instance = Instance(); + DevicePtr device; + auto config = PropertyObject(); + config.addProperty(IntProperty(PROPERTY_NAME_CLIENT_CONNECT_TIMEOUT, 1000)); + ASSERT_NO_THROW(device = instance.addDevice(MqttConnectionString(), config)); + ASSERT_EQ(device.getStatusContainer().getStatus("ComponentStatus"), + Enumeration("ComponentStatusType", "Ok", instance.getContext().getTypeManager())); + ASSERT_EQ(device.getPropertyValue(PROPERTY_NAME_CLIENT_CONNECT_TIMEOUT), 1000); +} + +TEST_F(MqttDeviceTest, CreatingSeveralMqttDevices) +{ + const auto instance = Instance(); + DevicePtr device; + ASSERT_NO_THROW(device = instance.addDevice(MqttConnectionString("127.0.0.1", 1883))); + ASSERT_EQ(device.getStatusContainer().getStatus("ComponentStatus"), + Enumeration("ComponentStatusType", "Ok", instance.getContext().getTypeManager())); + DevicePtr anotherDevice; + ASSERT_NO_THROW(anotherDevice = instance.addDevice(MqttConnectionString("127.0.0.1", 1883))); + ASSERT_EQ(anotherDevice.getStatusContainer().getStatus("ComponentStatus"), + Enumeration("ComponentStatusType", "Ok", instance.getContext().getTypeManager())); + ASSERT_EQ(instance.getDevices().getCount(), 2u); +} + +TEST_F(MqttDeviceTest, RemovingMqttDevice) +{ + const auto instance = Instance(); + DevicePtr device; + ASSERT_NO_THROW(device = instance.addDevice(MqttConnectionString(), DaqMqttDeviceConfig())); + ASSERT_EQ(instance.getDevices().getCount(), 1u); + ASSERT_NO_THROW(instance.removeDevice(device)); + ASSERT_EQ(instance.getDevices().getCount(), 0u); +} + +TEST_F(MqttDeviceTest, CheckMqttDeviceFunctionalBlocks) +{ + StartUp(); + DictPtr fbTypes; + ASSERT_NO_THROW(fbTypes = mqttDevice.getAvailableFunctionBlockTypes()); + ASSERT_EQ(fbTypes.getCount(), 2u); + ASSERT_TRUE(fbTypes.hasKey(SUB_FB_NAME)); + ASSERT_TRUE(fbTypes.hasKey(PUB_FB_NAME)); +} + +// --------------------------------------------------------------------------- +// Connection string handling +// --------------------------------------------------------------------------- + +TEST_F(MqttDeviceTest, ConnectionStringWithPort) +{ + const auto instance = Instance(); + DevicePtr device; + ASSERT_NO_THROW(device = instance.addDevice("daq.mqtt://127.0.0.1:1883")); + + ASSERT_EQ(device.getPropertyValue(PROPERTY_NAME_CLIENT_BROKER_PORT), 1883); + ASSERT_EQ(device.getInfo().getConnectionString(), "daq.mqtt://127.0.0.1:1883"); +} + +TEST_F(MqttDeviceTest, ConnectionStringWithoutPortFallsBackToProperty) +{ + const auto instance = Instance(); + auto config = DaqMqttDeviceConfig(); + config.setPropertyValue(PROPERTY_NAME_CLIENT_BROKER_PORT, 1883); + + DevicePtr device; + ASSERT_NO_THROW(device = instance.addDevice("daq.mqtt://127.0.0.1", config)); + + ASSERT_EQ(device.getPropertyValue(PROPERTY_NAME_CLIENT_BROKER_PORT), 1883); + // The normalised connection string always carries the effective port. + ASSERT_EQ(device.getInfo().getConnectionString(), "daq.mqtt://127.0.0.1:1883"); +} + +TEST_F(MqttDeviceTest, ConnectionStringPortBeatsProperty) +{ + const auto instance = Instance(); + auto config = DaqMqttDeviceConfig(); + // Points at a port nothing listens on; the connection string must win, so the device connects. + config.setPropertyValue(PROPERTY_NAME_CLIENT_BROKER_PORT, UNREACHABLE_PORT); + + DevicePtr device; + ASSERT_NO_THROW(device = instance.addDevice("daq.mqtt://127.0.0.1:1883", config)); + + // The property reflects the port actually used, not the one that came in through the config. + ASSERT_EQ(device.getPropertyValue(PROPERTY_NAME_CLIENT_BROKER_PORT), 1883); + ASSERT_EQ(device.getInfo().getConnectionString(), "daq.mqtt://127.0.0.1:1883"); +} + +TEST_F(MqttDeviceTest, UnknownConnectionStringPrefixThrows) +{ + const auto instance = Instance(); + ASSERT_ANY_THROW(instance.addDevice("daq.nosuchproto://127.0.0.1:1883")); +} + +// --------------------------------------------------------------------------- +// Connection string parsing +// --------------------------------------------------------------------------- + +TEST_F(MqttDeviceTest, ParseConnectionStringIpv4) +{ + const auto address = MqttStreamingModule::parseConnectionString("daq.mqtt://127.0.0.1:1883"); + ASSERT_EQ(address.host, "127.0.0.1"); + ASSERT_EQ(address.port, 1883); +} + +TEST_F(MqttDeviceTest, ParseConnectionStringHostname) +{ + const auto address = MqttStreamingModule::parseConnectionString("daq.mqtt://broker.emqx.io:1883"); + ASSERT_EQ(address.host, "broker.emqx.io"); + ASSERT_EQ(address.port, 1883); +} + +TEST_F(MqttDeviceTest, ParseConnectionStringWithoutPort) +{ + // Port 0 means "not given"; onCreateDevice then falls back to the Port property. + const auto address = MqttStreamingModule::parseConnectionString("daq.mqtt://127.0.0.1"); + ASSERT_EQ(address.host, "127.0.0.1"); + ASSERT_EQ(address.port, 0); +} + +TEST_F(MqttDeviceTest, ParseConnectionStringIpv6) +{ + const auto address = MqttStreamingModule::parseConnectionString("daq.mqtt://[::1]:1883"); + // The brackets stay part of the host - the MQTT client URL needs them. + ASSERT_EQ(address.host, "[::1]"); + ASSERT_EQ(address.port, 1883); +} + +TEST_F(MqttDeviceTest, ParseConnectionStringIpv6WithoutPort) +{ + const auto address = MqttStreamingModule::parseConnectionString("daq.mqtt://[fe80::1]"); + ASSERT_EQ(address.host, "[fe80::1]"); + ASSERT_EQ(address.port, 0); +} + +TEST_F(MqttDeviceTest, ParseConnectionStringRejectsForeignPrefix) +{ + ASSERT_THROW(MqttStreamingModule::parseConnectionString("daq.opcua://127.0.0.1:1883"), InvalidParameterException); +} + +TEST_F(MqttDeviceTest, ParseConnectionStringRejectsMissingPrefix) +{ + ASSERT_THROW(MqttStreamingModule::parseConnectionString("127.0.0.1:1883"), InvalidParameterException); +} + +TEST_F(MqttDeviceTest, ParseConnectionStringRejectsMissingHost) +{ + ASSERT_THROW(MqttStreamingModule::parseConnectionString("daq.mqtt://"), InvalidParameterException); +} + +TEST_F(MqttDeviceTest, ParseConnectionStringRejectsPortOutOfRange) +{ + ASSERT_THROW(MqttStreamingModule::parseConnectionString("daq.mqtt://127.0.0.1:0"), InvalidParameterException); + ASSERT_THROW(MqttStreamingModule::parseConnectionString("daq.mqtt://127.0.0.1:70000"), InvalidParameterException); +} + +TEST_F(MqttDeviceTest, FormatConnectionStringAlwaysCarriesPort) +{ + MqttStreamingModule::BrokerAddress address; + address.host = "127.0.0.1"; + address.port = 1883; + ASSERT_EQ(MqttStreamingModule::formatConnectionString(address), "daq.mqtt://127.0.0.1:1883"); + + address.host = "[::1]"; + ASSERT_EQ(MqttStreamingModule::formatConnectionString(address), "daq.mqtt://[::1]:1883"); +} + +TEST_F(MqttDeviceTest, UnreachableBrokerThrows) +{ + const auto instance = Instance(); + auto config = DaqMqttDeviceConfig(); + config.setPropertyValue(PROPERTY_NAME_CLIENT_CONNECT_TIMEOUT, 500); + + ASSERT_ANY_THROW(instance.addDevice(MqttConnectionString("127.0.0.1", UNREACHABLE_PORT), config)); + ASSERT_EQ(instance.getDevices().getCount(), 0u); +} + +// --------------------------------------------------------------------------- +// Device info and connection status +// --------------------------------------------------------------------------- + +TEST_F(MqttDeviceTest, DeviceInfoContent) +{ + StartUp(); + + const auto info = mqttDevice.getInfo(); + ASSERT_TRUE(info.assigned()); + ASSERT_EQ(info.getConnectionString(), MqttConnectionString()); + ASSERT_EQ(info.getName(), CLIENT_DEVICE_TYPE_NAME); + + const auto deviceType = info.getDeviceType(); + ASSERT_TRUE(deviceType.assigned()); + ASSERT_EQ(deviceType.getId(), CLIENT_DEVICE_TYPE_ID); + ASSERT_EQ(deviceType.getConnectionStringPrefix(), CLIENT_DEVICE_CONN_PREFIX); + + const auto connectionInfo = info.getConfigurationConnectionInfo(); + ASSERT_TRUE(connectionInfo.assigned()); + ASSERT_EQ(connectionInfo.getProtocolId(), CLIENT_DEVICE_TYPE_ID); + ASSERT_EQ(connectionInfo.getProtocolType(), ProtocolType::Unknown); + ASSERT_EQ(connectionInfo.getConnectionType(), "TCP/IP"); + ASSERT_EQ(connectionInfo.getPort(), DEFAULT_PORT); + ASSERT_EQ(connectionInfo.getPrefix(), CLIENT_DEVICE_CONN_PREFIX); + ASSERT_EQ(connectionInfo.getConnectionString(), MqttConnectionString()); + ASSERT_EQ(connectionInfo.getAddresses().getCount(), 1u); + ASSERT_EQ(connectionInfo.getAddresses()[0], DEFAULT_BROKER_ADDRESS); +} + +TEST_F(MqttDeviceTest, ConfigurationStatusConnected) +{ + StartUp(); + + const auto statuses = mqttDevice.getConnectionStatusContainer(); + ASSERT_TRUE(statuses.assigned()); + ASSERT_TRUE(statuses.getStatuses().hasKey("ConfigurationStatus")); + ASSERT_EQ(statuses.getStatus("ConfigurationStatus"), + Enumeration("ConnectionStatusType", "Connected", daqInstance.getContext().getTypeManager())); +} + +TEST_F(MqttDeviceTest, ConnectionStatusKeyedByConnectionString) +{ + StartUp(); + + // The container is keyed by connection string; the alias is what clients read. + const auto statuses = mqttDevice.getConnectionStatusContainer().getStatuses(); + ASSERT_EQ(statuses.getCount(), 1u); + ASSERT_TRUE(statuses.hasKey("ConfigurationStatus")); +} diff --git a/modules/mqtt_streaming_module/tests/test_mqtt_client_fb.cpp b/modules/mqtt_streaming_module/tests/test_mqtt_client_fb.cpp deleted file mode 100644 index cfb2ff0c..00000000 --- a/modules/mqtt_streaming_module/tests/test_mqtt_client_fb.cpp +++ /dev/null @@ -1,149 +0,0 @@ -#include "test_daq_test_helper.h" -#include -#include -#include -#include -#include - -using namespace daq; -using namespace daq::modules::mqtt_streaming_module; - -namespace daq::modules::mqtt_streaming_module -{ -class MqttFbTest : public testing::Test, public DaqTestHelper -{ -}; -} // namespace daq::modules::mqtt_streaming_module - -TEST_F(MqttFbTest, DefaultMqttFbConfig) -{ - const auto module = CreateModule(); - - DictPtr types; - ASSERT_NO_THROW(types = module.getAvailableFunctionBlockTypes()); - ASSERT_EQ(types.getCount(), 1u); - - ASSERT_TRUE(types.hasKey(CLIENT_FB_NAME)); - auto defaultConfig = types.get(CLIENT_FB_NAME).createDefaultConfig(); - ASSERT_TRUE(defaultConfig.assigned()); - - ASSERT_EQ(defaultConfig.getAllProperties().getCount(), 5u); - - ASSERT_TRUE(defaultConfig.hasProperty(PROPERTY_NAME_CLIENT_BROKER_ADDRESS)); - ASSERT_TRUE(defaultConfig.hasProperty(PROPERTY_NAME_CLIENT_BROKER_PORT)); - ASSERT_TRUE(defaultConfig.hasProperty(PROPERTY_NAME_CLIENT_USERNAME)); - ASSERT_TRUE(defaultConfig.hasProperty(PROPERTY_NAME_CLIENT_PASSWORD)); - ASSERT_TRUE(defaultConfig.hasProperty(PROPERTY_NAME_CLIENT_CONNECT_TIMEOUT)); - - ASSERT_EQ(defaultConfig.getProperty(PROPERTY_NAME_CLIENT_BROKER_ADDRESS).getValueType(), CoreType::ctString); - ASSERT_EQ(defaultConfig.getProperty(PROPERTY_NAME_CLIENT_BROKER_PORT).getValueType(), CoreType::ctInt); - ASSERT_EQ(defaultConfig.getProperty(PROPERTY_NAME_CLIENT_USERNAME).getValueType(), CoreType::ctString); - ASSERT_EQ(defaultConfig.getProperty(PROPERTY_NAME_CLIENT_PASSWORD).getValueType(), CoreType::ctString); - ASSERT_EQ(defaultConfig.getProperty(PROPERTY_NAME_CLIENT_CONNECT_TIMEOUT).getValueType(), CoreType::ctInt); - - ASSERT_EQ(defaultConfig.getPropertyValue(PROPERTY_NAME_CLIENT_BROKER_ADDRESS), DEFAULT_BROKER_ADDRESS); - ASSERT_EQ(defaultConfig.getPropertyValue(PROPERTY_NAME_CLIENT_BROKER_PORT), DEFAULT_PORT); - ASSERT_EQ(defaultConfig.getPropertyValue(PROPERTY_NAME_CLIENT_USERNAME), DEFAULT_USERNAME); - ASSERT_EQ(defaultConfig.getPropertyValue(PROPERTY_NAME_CLIENT_PASSWORD), DEFAULT_PASSWORD); - ASSERT_EQ(defaultConfig.getPropertyValue(PROPERTY_NAME_CLIENT_CONNECT_TIMEOUT), DEFAULT_INIT_TIMEOUT); -} - -TEST_F(MqttFbTest, MissingPasswordProperty) -{ - const auto instance = Instance(); - daq::FunctionBlockPtr fb; - ASSERT_NO_THROW(fb = instance.addFunctionBlock(CLIENT_FB_NAME)); - ASSERT_EQ(fb.getStatusContainer().getStatus("ComponentStatus"), - Enumeration("ComponentStatusType", "Ok", instance.getContext().getTypeManager())); - - ASSERT_FALSE(fb.hasProperty(PROPERTY_NAME_CLIENT_PASSWORD)); -} - -TEST_F(MqttFbTest, CreatingMqttFbWithDefaultConfig) -{ - const auto instance = Instance(); - daq::FunctionBlockPtr fb; - ASSERT_NO_THROW(fb = instance.addFunctionBlock(CLIENT_FB_NAME)); - ASSERT_EQ(fb.getStatusContainer().getStatus("ComponentStatus"), - Enumeration("ComponentStatusType", "Ok", instance.getContext().getTypeManager())); - - auto fbs = instance.getFunctionBlocks(); - bool contain = false; - daq::FunctionBlockPtr fbFromList; - for (const auto& fbInst : fbs) - { - contain = (fbInst.getName().toStdString().find(MQTT_LOCAL_CLIENT_FB_ID_PREFIX) != std::string::npos); - if (contain) - { - fbFromList = fbInst; - break; - } - } - ASSERT_TRUE(contain); - ASSERT_TRUE(fbFromList.assigned()); - ASSERT_TRUE(fbFromList == fb); -} - -TEST_F(MqttFbTest, CreatingMqttFbWithCustomConfig) -{ - const auto instance = Instance(); - daq::FunctionBlockPtr fb; - auto config = DaqMqttFbConfig(); - ASSERT_NO_THROW(fb = instance.addFunctionBlock(CLIENT_FB_NAME, config)); - ASSERT_EQ(fb.getStatusContainer().getStatus("ComponentStatus"), - Enumeration("ComponentStatusType", "Ok", instance.getContext().getTypeManager())); -} - -TEST_F(MqttFbTest, CreatingMqttFbWithEmptyConfig) -{ - const auto instance = Instance(); - daq::FunctionBlockPtr fb; - auto config = PropertyObject(); - ASSERT_NO_THROW(fb = instance.addFunctionBlock(CLIENT_FB_NAME, config)); - ASSERT_EQ(fb.getStatusContainer().getStatus("ComponentStatus"), - Enumeration("ComponentStatusType", "Ok", instance.getContext().getTypeManager())); -} - -TEST_F(MqttFbTest, CreatingMqttFbWithPartialConfig) -{ - const auto instance = Instance(); - daq::FunctionBlockPtr fb; - auto config = PropertyObject(); - config.addProperty(IntProperty(PROPERTY_NAME_CLIENT_CONNECT_TIMEOUT, 1000)); - ASSERT_NO_THROW(fb = instance.addFunctionBlock(CLIENT_FB_NAME, config)); - ASSERT_EQ(fb.getStatusContainer().getStatus("ComponentStatus"), - Enumeration("ComponentStatusType", "Ok", instance.getContext().getTypeManager())); -} - -TEST_F(MqttFbTest, DISABLED_CreatingSeveralMqttFbs) -{ - const auto instance = Instance(); - daq::FunctionBlockPtr fb; - ASSERT_NO_THROW(fb = instance.addFunctionBlock(CLIENT_FB_NAME, DaqMqttFbConfig("127.0.0.1", 1883))); - ASSERT_EQ(fb.getStatusContainer().getStatus("ComponentStatus"), - Enumeration("ComponentStatusType", "Ok", instance.getContext().getTypeManager())); - daq::FunctionBlockPtr anotherFb; - ASSERT_NO_THROW(anotherFb = instance.addFunctionBlock(CLIENT_FB_NAME, DaqMqttFbConfig("127.0.0.1", 1884))); - ASSERT_EQ(anotherFb.getStatusContainer().getStatus("ComponentStatus"), - Enumeration("ComponentStatusType", "Ok", instance.getContext().getTypeManager())); - ASSERT_EQ(instance.getFunctionBlocks().getCount(), 2u); -} - -TEST_F(MqttFbTest, RemovingMqttFb) -{ - const auto instance = Instance(); - daq::FunctionBlockPtr fb; - auto config = DaqMqttFbConfig(); - ASSERT_NO_THROW(fb = instance.addFunctionBlock(CLIENT_FB_NAME, config)); - ASSERT_NO_THROW(instance.removeFunctionBlock(fb)); -} - -TEST_F(MqttFbTest, CheckMqttFbFunctionalBlocks) -{ - StartUp(); - daq::DictPtr fbTypes; - ASSERT_NO_THROW(fbTypes = clientMqttFb.getAvailableFunctionBlockTypes()); - ASSERT_GE(fbTypes.getCount(), 2); - ASSERT_TRUE(fbTypes.hasKey(SUB_FB_NAME)); - ASSERT_TRUE(fbTypes.hasKey(PUB_FB_NAME)); -} diff --git a/modules/mqtt_streaming_module/tests/test_mqtt_json_decoder_fb.cpp b/modules/mqtt_streaming_module/tests/test_mqtt_json_decoder_fb.cpp index 513902af..b88abb2c 100644 --- a/modules/mqtt_streaming_module/tests/test_mqtt_json_decoder_fb.cpp +++ b/modules/mqtt_streaming_module/tests/test_mqtt_json_decoder_fb.cpp @@ -1059,7 +1059,7 @@ TEST_F(MqttJsonDecoderFbTest, DataTransferSeveralSignals) const auto topic = buildTopicName(); DaqInstanceInit(); - auto clientFb0 = DaqAddClientMqttFb("127.0.0.1", DEFAULT_PORT); + auto clientFb0 = DaqAddMqttDevice("127.0.0.1", DEFAULT_PORT); auto jsonFb0 = AddSubFb(topic); auto decoderFb0 = AddDecoderFb(valueF0, DDSM::ExtractFromMessage, tsF); auto decoderFb1 = AddDecoderFb(valueF1, DDSM::ExtractFromMessage, tsF); @@ -1146,7 +1146,7 @@ TEST_F(MqttJsonDecoderFbTest, DataTransferMissingFieldSeveralSignals) const auto topic = buildTopicName(); DaqInstanceInit(); - auto clientFb0 = DaqAddClientMqttFb("127.0.0.1", DEFAULT_PORT); + auto clientFb0 = DaqAddMqttDevice("127.0.0.1", DEFAULT_PORT); auto jsonFb0 = AddSubFb(topic); auto decoderFb0 = AddDecoderFb(valueF0, DDSM::ExtractFromMessage, tsF); auto decoderFb1 = AddDecoderFb(valueF1, DDSM::ExtractFromMessage, tsF); @@ -1227,11 +1227,11 @@ TEST_F(MqttJsonFbCommunicationTest, DISABLED_FullDataTransferFor2MqttFbs) const std::string topic1 = buildTopicName("1"); DaqInstanceInit(); - auto clientFb0 = DaqAddClientMqttFb("127.0.0.1", 1883); + auto clientFb0 = DaqAddMqttDevice("127.0.0.1", 1883); auto jsonFb0 = AddSubFb(topic0); auto decoderFb0 = AddDecoderFb(valueF, DDSM::ExtractFromMessage, tsF); - auto clientFb1 = DaqAddClientMqttFb("127.0.0.1", 1884); + auto clientFb1 = DaqAddMqttDevice("127.0.0.1", 1884); auto jsonFb1 = AddSubFb(topic1); auto decoderFb1 = AddDecoderFb(valueF, DDSM::ExtractFromMessage, tsF); diff --git a/modules/mqtt_streaming_module/tests/test_mqtt_publisher_fb.cpp b/modules/mqtt_streaming_module/tests/test_mqtt_publisher_fb.cpp index b808626d..1f4aa2d0 100644 --- a/modules/mqtt_streaming_module/tests/test_mqtt_publisher_fb.cpp +++ b/modules/mqtt_streaming_module/tests/test_mqtt_publisher_fb.cpp @@ -237,10 +237,10 @@ class MqttPublisherFbHelper : public DaqTestHelper void CreatePublisherFB() { - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_PUB_QOS, 2); config.setPropertyValue(PROPERTY_NAME_PUB_PREVIEW_SIGNAL, True); - fb = clientMqttFb.addFunctionBlock(PUB_FB_NAME, config); + fb = mqttDevice.addFunctionBlock(PUB_FB_NAME, config); } void CreatePublisherFB(bool multiTopic, @@ -251,7 +251,7 @@ class MqttPublisherFbHelper : public DaqTestHelper int qos = 2, uint32_t readPeriod = 20) { - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_PUB_MODE, 0); config.setPropertyValue(PROPERTY_NAME_PUB_TOPIC_MODE, multiTopic ? 1 : 0); config.setPropertyValue(PROPERTY_NAME_PUB_GROUP_VALUES, groupV ? True : False); @@ -261,18 +261,18 @@ class MqttPublisherFbHelper : public DaqTestHelper config.setPropertyValue(PROPERTY_NAME_PUB_READ_PERIOD, readPeriod); config.setPropertyValue(PROPERTY_NAME_PUB_TOPIC_NAME, topicName); config.setPropertyValue(PROPERTY_NAME_PUB_PREVIEW_SIGNAL, True); - fb = clientMqttFb.addFunctionBlock(PUB_FB_NAME, config); + fb = mqttDevice.addFunctionBlock(PUB_FB_NAME, config); } void CreateRawPublisherFB(int qos = 2, uint32_t readPeriod = 20) { - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_PUB_MODE, 1); config.setPropertyValue(PROPERTY_NAME_PUB_QOS, qos); config.setPropertyValue(PROPERTY_NAME_PUB_READ_PERIOD, readPeriod); config.setPropertyValue(PROPERTY_NAME_PUB_PREVIEW_SIGNAL, False); - fb = clientMqttFb.addFunctionBlock(PUB_FB_NAME, config); + fb = mqttDevice.addFunctionBlock(PUB_FB_NAME, config); } bool CreateSubscriber(std::string postfix = "_subscriberId") @@ -622,7 +622,7 @@ TEST_F(MqttPublisherFbTest, DefaultConfig) daq::DictPtr fbTypes; daq::FunctionBlockTypePtr fbt; daq::PropertyObjectPtr defaultConfig; - ASSERT_NO_THROW(fbTypes = clientMqttFb.getAvailableFunctionBlockTypes()); + ASSERT_NO_THROW(fbTypes = mqttDevice.getAvailableFunctionBlockTypes()); ASSERT_NO_THROW(fbt = fbTypes.get(PUB_FB_NAME)); ASSERT_NO_THROW(defaultConfig = fbt.createDefaultConfig()); @@ -723,7 +723,7 @@ TEST_F(MqttPublisherFbTest, PropertyVisibility) TEST_F(MqttPublisherFbTest, Config) { StartUp(); - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_PUB_MODE, 0); config.setPropertyValue(PROPERTY_NAME_PUB_TOPIC_MODE, 1); @@ -735,7 +735,7 @@ TEST_F(MqttPublisherFbTest, Config) config.setPropertyValue(PROPERTY_NAME_PUB_TOPIC_NAME, buildTopicName()); config.setPropertyValue(PROPERTY_NAME_PUB_PREVIEW_SIGNAL, True); daq::FunctionBlockPtr fb; - ASSERT_NO_THROW(fb = clientMqttFb.addFunctionBlock(PUB_FB_NAME, config)); + ASSERT_NO_THROW(fb = mqttDevice.addFunctionBlock(PUB_FB_NAME, config)); ASSERT_EQ(fb.getStatusContainer().getStatus("ComponentStatus"), Enumeration("ComponentStatusType", "Warning", daqInstance.getContext().getTypeManager())); SignalHelper helper{}; @@ -770,7 +770,7 @@ TEST_F(MqttPublisherFbTest, Creation) { StartUp(); daq::FunctionBlockPtr fb; - ASSERT_NO_THROW(fb = clientMqttFb.addFunctionBlock(PUB_FB_NAME)); + ASSERT_NO_THROW(fb = mqttDevice.addFunctionBlock(PUB_FB_NAME)); SignalHelper helper{}; fb.getInputPorts()[0].connect(helper.signal0); ASSERT_EQ(fb.getStatusContainer().getStatus("ComponentStatus"), @@ -787,7 +787,7 @@ TEST_F(MqttPublisherFbTest, TwoFbCreation) SignalHelper helper{}; { daq::FunctionBlockPtr fb; - ASSERT_NO_THROW(fb = clientMqttFb.addFunctionBlock(PUB_FB_NAME)); + ASSERT_NO_THROW(fb = mqttDevice.addFunctionBlock(PUB_FB_NAME)); fb.getInputPorts()[0].connect(helper.signal0); ASSERT_EQ(fb.getStatusContainer().getStatus("ComponentStatus"), Enumeration("ComponentStatusType", "Ok", daqInstance.getContext().getTypeManager())); @@ -798,7 +798,7 @@ TEST_F(MqttPublisherFbTest, TwoFbCreation) } { daq::FunctionBlockPtr fb; - ASSERT_NO_THROW(fb = clientMqttFb.addFunctionBlock(PUB_FB_NAME)); + ASSERT_NO_THROW(fb = mqttDevice.addFunctionBlock(PUB_FB_NAME)); fb.getInputPorts()[0].connect(helper.signal0); ASSERT_EQ(fb.getStatusContainer().getStatus("ComponentStatus"), Enumeration("ComponentStatusType", "Ok", daqInstance.getContext().getTypeManager())); @@ -807,7 +807,7 @@ TEST_F(MqttPublisherFbTest, TwoFbCreation) static_cast(MqttPublisherFbImpl::SettingStatus::Valid), daqInstance.getContext().getTypeManager())); } - auto fbs = clientMqttFb.getFunctionBlocks(); + auto fbs = mqttDevice.getFunctionBlocks(); ASSERT_EQ(fbs.getCount(), 2u); } @@ -815,7 +815,7 @@ TEST_F(MqttPublisherFbTest, CreationWithDefaultConfig) { StartUp(); daq::FunctionBlockPtr fb; - ASSERT_NO_THROW(fb = clientMqttFb.addFunctionBlock(PUB_FB_NAME)); + ASSERT_NO_THROW(fb = mqttDevice.addFunctionBlock(PUB_FB_NAME)); auto signals = fb.getSignals(); ASSERT_EQ(signals.getCount(), 0u); ASSERT_EQ(fb.getStatusContainer().getStatus("ComponentStatus"), @@ -836,7 +836,7 @@ TEST_F(MqttPublisherFbTest, CreationWithPartialConfig) daq::FunctionBlockPtr fb; auto config = PropertyObject(); config.addProperty(IntProperty(PROPERTY_NAME_PUB_READ_PERIOD, 20)); - ASSERT_NO_THROW(fb = clientMqttFb.addFunctionBlock(PUB_FB_NAME, config)); + ASSERT_NO_THROW(fb = mqttDevice.addFunctionBlock(PUB_FB_NAME, config)); SignalHelper helper{}; fb.getInputPorts()[0].connect(helper.signal0); ASSERT_EQ(fb.getStatusContainer().getStatus("ComponentStatus"), @@ -853,11 +853,11 @@ TEST_F(MqttPublisherFbTest, ConnectToPort) StatusHelper::addTypesToTypeManager(MQTT_PUB_FB_SIG_STATUS_TYPE, MQTT_PUB_FB_SIG_STATUS_NAME, MqttPublisherFbImpl::signalStatusMap, - clientMqttFb.getContext().getTypeManager()); + mqttDevice.getContext().getTypeManager()); StatusHelper::addTypesToTypeManager(MQTT_PUB_FB_PUB_STATUS_TYPE, MQTT_PUB_FB_PUB_STATUS_NAME, MqttPublisherFbImpl::publishingStatusMap, - clientMqttFb.getContext().getTypeManager()); + mqttDevice.getContext().getTypeManager()); const auto sigStValid = EnumerationWithIntValue(MQTT_PUB_FB_SIG_STATUS_TYPE, static_cast(MqttPublisherFbImpl::SignalStatus::Valid), daqInstance.getContext().getTypeManager()); @@ -873,7 +873,7 @@ TEST_F(MqttPublisherFbTest, ConnectToPort) { daq::FunctionBlockPtr fb; - ASSERT_NO_THROW(fb = clientMqttFb.addFunctionBlock(PUB_FB_NAME)); + ASSERT_NO_THROW(fb = mqttDevice.addFunctionBlock(PUB_FB_NAME)); ASSERT_EQ(fb.getStatusContainer().getStatus(MQTT_PUB_FB_SIG_STATUS_NAME), sigStNotConnected); auto help = SignalHelper(); @@ -920,10 +920,10 @@ TEST_F(MqttPublisherFbTest, ConnectToPort) { daq::FunctionBlockPtr fb; - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_PUB_TOPIC_MODE, 1); config.setPropertyValue(PROPERTY_NAME_PUB_TOPIC_NAME, String(buildTopicName())); - ASSERT_NO_THROW(fb = clientMqttFb.addFunctionBlock(PUB_FB_NAME, config)); + ASSERT_NO_THROW(fb = mqttDevice.addFunctionBlock(PUB_FB_NAME, config)); auto help = SignalHelper(); auto signal0 = help.createSignal(DataDescriptorBuilder().setRule(LinearDataRule(2, 3)).setTickResolution(Ratio(1, 1000))); @@ -937,10 +937,10 @@ TEST_F(MqttPublisherFbTest, ConnectToPort) { daq::FunctionBlockPtr fb; - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_PUB_TOPIC_MODE, 1); config.setPropertyValue(PROPERTY_NAME_PUB_TOPIC_NAME, String(buildTopicName())); - ASSERT_NO_THROW(fb = clientMqttFb.addFunctionBlock(PUB_FB_NAME, config)); + ASSERT_NO_THROW(fb = mqttDevice.addFunctionBlock(PUB_FB_NAME, config)); auto help = SignalHelper(); auto signal0 = help.createSignal(DataDescriptorBuilder().setRule(LinearDataRule(1, 3)).setTickResolution(Ratio(1, 1000))); @@ -957,10 +957,10 @@ TEST_F(MqttPublisherFbTest, ConnectToPort) { daq::FunctionBlockPtr fb; - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_PUB_TOPIC_MODE, 1); config.setPropertyValue(PROPERTY_NAME_PUB_TOPIC_NAME, String(buildTopicName())); - ASSERT_NO_THROW(fb = clientMqttFb.addFunctionBlock(PUB_FB_NAME, config)); + ASSERT_NO_THROW(fb = mqttDevice.addFunctionBlock(PUB_FB_NAME, config)); auto help = SignalHelper(); auto signal0 = help.createSignal(DataDescriptorBuilder().setRule(LinearDataRule(2, 3)).setTickResolution(Ratio(1, 1000))); @@ -977,10 +977,10 @@ TEST_F(MqttPublisherFbTest, ConnectToPort) { daq::FunctionBlockPtr fb; - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_PUB_TOPIC_MODE, 1); config.setPropertyValue(PROPERTY_NAME_PUB_TOPIC_NAME, String(buildTopicName())); - ASSERT_NO_THROW(fb = clientMqttFb.addFunctionBlock(PUB_FB_NAME, config)); + ASSERT_NO_THROW(fb = mqttDevice.addFunctionBlock(PUB_FB_NAME, config)); auto help = SignalHelper(); auto signal0 = help.createSignal(DataDescriptorBuilder().setRule(LinearDataRule(1, 3)).setTickResolution(Ratio(1, 500))); @@ -1006,10 +1006,10 @@ TEST_F(MqttPublisherFbTest, PreviewSignals) StartUp(); daq::FunctionBlockPtr fb; - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_PUB_TOPIC_MODE, 0); config.setPropertyValue(PROPERTY_NAME_PUB_PREVIEW_SIGNAL, False); - ASSERT_NO_THROW(fb = clientMqttFb.addFunctionBlock(PUB_FB_NAME, config)); + ASSERT_NO_THROW(fb = mqttDevice.addFunctionBlock(PUB_FB_NAME, config)); auto help = SignalHelper(); ASSERT_EQ(fb.getSignals().getCount(), 0u); @@ -1068,9 +1068,9 @@ TEST_F(MqttPublisherFbTest, TopicsList) { daq::FunctionBlockPtr fb; - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_PUB_TOPIC_MODE, 0); - ASSERT_NO_THROW(fb = clientMqttFb.addFunctionBlock(PUB_FB_NAME, config)); + ASSERT_NO_THROW(fb = mqttDevice.addFunctionBlock(PUB_FB_NAME, config)); auto help = SignalHelper(); ASSERT_NO_THROW(fb.getInputPorts()[0].connect(help.signal0)); @@ -1082,9 +1082,9 @@ TEST_F(MqttPublisherFbTest, TopicsList) { daq::FunctionBlockPtr fb; - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_PUB_TOPIC_MODE, 1); - ASSERT_NO_THROW(fb = clientMqttFb.addFunctionBlock(PUB_FB_NAME, config)); + ASSERT_NO_THROW(fb = mqttDevice.addFunctionBlock(PUB_FB_NAME, config)); auto help = SignalHelper(); ASSERT_NO_THROW(fb.getInputPorts()[0].connect(help.signal0)); @@ -1099,11 +1099,11 @@ TEST_F(MqttPublisherFbTest, WrongConfig) { StartUp(); daq::FunctionBlockPtr fb; - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(PUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_PUB_MODE, 0); config.setPropertyValue(PROPERTY_NAME_PUB_TOPIC_MODE, 1); config.setPropertyValue(PROPERTY_NAME_PUB_TOPIC_NAME, String("/test/#")); - ASSERT_NO_THROW(fb = clientMqttFb.addFunctionBlock(PUB_FB_NAME, config)); + ASSERT_NO_THROW(fb = mqttDevice.addFunctionBlock(PUB_FB_NAME, config)); ASSERT_EQ(fb.getStatusContainer().getStatus("ComponentStatus"), Enumeration("ComponentStatusType", "Error", daqInstance.getContext().getTypeManager())); ASSERT_EQ(fb.getStatusContainer().getStatus(MQTT_PUB_FB_SET_STATUS_NAME), diff --git a/modules/mqtt_streaming_module/tests/test_mqtt_streaming_module.cpp b/modules/mqtt_streaming_module/tests/test_mqtt_streaming_module.cpp index bfc7cc20..917f9378 100644 --- a/modules/mqtt_streaming_module/tests/test_mqtt_streaming_module.cpp +++ b/modules/mqtt_streaming_module/tests/test_mqtt_streaming_module.cpp @@ -49,28 +49,30 @@ TEST_F(MqttStreamingClientModuleTest, VersionCorrect) ASSERT_EQ(version.getPatch(), MQTT_STREAM_MODULE_PATCH_VERSION); } -TEST_F(MqttStreamingClientModuleTest, MqttFbAvailable) +TEST_F(MqttStreamingClientModuleTest, MqttDeviceAvailable) { auto module = CreateModule(); - DictPtr fbt; - ASSERT_NO_THROW(fbt = module.getAvailableFunctionBlockTypes()); - ASSERT_EQ(fbt.getCount(), 1u); + DictPtr deviceTypes; + ASSERT_NO_THROW(deviceTypes = module.getAvailableDeviceTypes()); + ASSERT_EQ(deviceTypes.getCount(), 1u); } TEST_F(MqttStreamingClientModuleTest, GetAvailableComponentTypes) { const auto module = CreateModule(); + // The MQTT client is a device now; the module offers no top-level function block type. DictPtr functionBlockTypes; ASSERT_NO_THROW(functionBlockTypes = module.getAvailableFunctionBlockTypes()); - ASSERT_EQ(functionBlockTypes.getCount(), 1u); - ASSERT_TRUE(functionBlockTypes.hasKey(CLIENT_FB_NAME)); - ASSERT_EQ(functionBlockTypes.get(CLIENT_FB_NAME).getId(), CLIENT_FB_NAME); + ASSERT_EQ(functionBlockTypes.getCount(), 0u); DictPtr deviceTypes; ASSERT_NO_THROW(deviceTypes = module.getAvailableDeviceTypes()); - ASSERT_EQ(deviceTypes.getCount(), 0u); + ASSERT_EQ(deviceTypes.getCount(), 1u); + ASSERT_TRUE(deviceTypes.hasKey(CLIENT_DEVICE_TYPE_ID)); + ASSERT_EQ(deviceTypes.get(CLIENT_DEVICE_TYPE_ID).getId(), CLIENT_DEVICE_TYPE_ID); + ASSERT_EQ(deviceTypes.get(CLIENT_DEVICE_TYPE_ID).getConnectionStringPrefix(), CLIENT_DEVICE_CONN_PREFIX); DictPtr serverTypes; @@ -94,11 +96,11 @@ TEST_F(MqttStreamingClientModuleTest, GetAvailableComponentTypes) ASSERT_EQ(versionInfoModule.getPatch(), MQTT_STREAM_MODULE_PATCH_VERSION); } - // Check module and version info for fb types - for (const auto& fbt : functionBlockTypes) + // Check module and version info for device types + for (const auto& devType : deviceTypes) { ModuleInfoPtr moduleInfo; - ASSERT_NO_THROW(moduleInfo = fbt.second.getModuleInfo()); + ASSERT_NO_THROW(moduleInfo = devType.second.getModuleInfo()); ASSERT_NE(moduleInfo, nullptr); ASSERT_EQ(moduleInfo.getName(), MODULE_NAME); ASSERT_EQ(moduleInfo.getId(), MODULE_ID); diff --git a/modules/mqtt_streaming_module/tests/test_mqtt_subscriber_fb.cpp b/modules/mqtt_streaming_module/tests/test_mqtt_subscriber_fb.cpp index 910bd5d3..11a90e08 100644 --- a/modules/mqtt_streaming_module/tests/test_mqtt_subscriber_fb.cpp +++ b/modules/mqtt_streaming_module/tests/test_mqtt_subscriber_fb.cpp @@ -202,7 +202,7 @@ TEST_F(MqttSubscriberFbTest, Config) config.setPropertyValue(PROPERTY_NAME_SUB_TOPIC, buildTopicName()); daq::FunctionBlockPtr subFb; - ASSERT_NO_THROW(subFb = clientMqttFb.addFunctionBlock(SUB_FB_NAME, config)); + ASSERT_NO_THROW(subFb = mqttDevice.addFunctionBlock(SUB_FB_NAME, config)); const auto allProperties = subFb.getAllProperties(); ASSERT_EQ(allProperties.getCount(), config.getAllProperties().getCount()); @@ -219,7 +219,7 @@ TEST_F(MqttSubscriberFbTest, CreationWithDefaultConfig) { StartUp(); daq::FunctionBlockPtr subFb; - ASSERT_NO_THROW(subFb = clientMqttFb.addFunctionBlock(SUB_FB_NAME)); + ASSERT_NO_THROW(subFb = mqttDevice.addFunctionBlock(SUB_FB_NAME)); EXPECT_EQ(subFb.getSignals(daq::search::Any()).getCount(), 0u); ASSERT_TRUE(waitStatusChange(1500, subFb, Enumeration("ComponentStatusType", "Error", daqInstance.getContext().getTypeManager()))); } @@ -231,7 +231,7 @@ TEST_F(MqttSubscriberFbTest, CreationWithPartialConfig) daq::FunctionBlockPtr subFb; auto config = PropertyObject(); config.addProperty(StringProperty(PROPERTY_NAME_SUB_TOPIC, String(buildTopicName()))); - ASSERT_NO_THROW(subFb = clientMqttFb.addFunctionBlock(SUB_FB_NAME, config)); + ASSERT_NO_THROW(subFb = mqttDevice.addFunctionBlock(SUB_FB_NAME, config)); EXPECT_EQ(subFb.getSignals(daq::search::Any()).getCount(), 0u); ASSERT_TRUE(waitStatusChange(1500, subFb, Enumeration("ComponentStatusType", "Warning", daqInstance.getContext().getTypeManager()))); } @@ -241,11 +241,11 @@ TEST_F(MqttSubscriberFbTest, CreationWithCustomConfig) // If FB has only one property, partial config is equivalent to custom config StartUp(); daq::FunctionBlockPtr subFb; - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_SUB_PREVIEW_SIGNAL, True); config.setPropertyValue(PROPERTY_NAME_SUB_TOPIC, buildTopicName()); config.setPropertyValue(PROPERTY_NAME_SUB_PREVIEW_SIGNAL_TS_MODE, static_cast(SDSM::SystemTime)); - ASSERT_NO_THROW(subFb = clientMqttFb.addFunctionBlock(SUB_FB_NAME, config)); + ASSERT_NO_THROW(subFb = mqttDevice.addFunctionBlock(SUB_FB_NAME, config)); EXPECT_EQ(subFb.getSignals(daq::search::Any()).getCount(), 2u); ASSERT_TRUE(waitStatusChange(1500, subFb, Enumeration("ComponentStatusType", "Warning", daqInstance.getContext().getTypeManager()))); } @@ -254,12 +254,12 @@ TEST_F(MqttSubscriberFbTest, PreviewSignal) { StartUp(); daq::FunctionBlockPtr subFb; - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_SUB_PREVIEW_SIGNAL, True); config.setPropertyValue(PROPERTY_NAME_SUB_PREVIEW_SIGNAL_IS_STRING, False); config.setPropertyValue(PROPERTY_NAME_SUB_TOPIC, buildTopicName()); config.setPropertyValue(PROPERTY_NAME_SUB_PREVIEW_SIGNAL_TS_MODE, static_cast(SDSM::SystemTime)); - ASSERT_NO_THROW(subFb = clientMqttFb.addFunctionBlock(SUB_FB_NAME, config)); + ASSERT_NO_THROW(subFb = mqttDevice.addFunctionBlock(SUB_FB_NAME, config)); ASSERT_EQ(subFb.getSignals().getCount(), 1u); EXPECT_EQ(subFb.getSignals()[0].getDescriptor().getSampleType(), daq::SampleType::Binary); subFb.setPropertyValue(PROPERTY_NAME_SUB_PREVIEW_SIGNAL_IS_STRING, True); @@ -270,12 +270,12 @@ TEST_F(MqttSubscriberFbTest, DomainForPreviewSignal) { StartUp(); daq::FunctionBlockPtr subFb; - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_SUB_PREVIEW_SIGNAL, True); config.setPropertyValue(PROPERTY_NAME_SUB_PREVIEW_SIGNAL_IS_STRING, False); config.setPropertyValue(PROPERTY_NAME_SUB_TOPIC, buildTopicName()); config.setPropertyValue(PROPERTY_NAME_SUB_PREVIEW_SIGNAL_TS_MODE, static_cast(SDSM::None)); - ASSERT_NO_THROW(subFb = clientMqttFb.addFunctionBlock(SUB_FB_NAME, config)); + ASSERT_NO_THROW(subFb = mqttDevice.addFunctionBlock(SUB_FB_NAME, config)); ASSERT_EQ(subFb.getSignals(daq::search::Any()).getCount(), 1u); subFb.setPropertyValue(PROPERTY_NAME_SUB_PREVIEW_SIGNAL_TS_MODE, static_cast(SDSM::SystemTime)); ASSERT_EQ(subFb.getSignals(daq::search::Any()).getCount(), 2u); @@ -289,11 +289,11 @@ TEST_F(MqttSubscriberFbTest, SubscriptionStatusWaitingForData) { StartUp(); - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_SUB_TOPIC, buildTopicName()); daq::FunctionBlockPtr subFb; - ASSERT_NO_THROW(subFb = clientMqttFb.addFunctionBlock(SUB_FB_NAME, config)); + ASSERT_NO_THROW(subFb = mqttDevice.addFunctionBlock(SUB_FB_NAME, config)); ASSERT_TRUE(waitStatusChange(1500, subFb, Enumeration("ComponentStatusType", "Warning", daqInstance.getContext().getTypeManager()))); } @@ -302,12 +302,12 @@ TEST_P(MqttSubscriberFbTopicPTest, CheckSubscriberFbTopic) auto [topic, result] = GetParam(); StartUp(); - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_SUB_TOPIC, topic); config.setPropertyValue(PROPERTY_NAME_SUB_PREVIEW_SIGNAL, True); daq::FunctionBlockPtr fb; - ASSERT_NO_THROW(fb = clientMqttFb.addFunctionBlock(SUB_FB_NAME, config)); + ASSERT_NO_THROW(fb = mqttDevice.addFunctionBlock(SUB_FB_NAME, config)); auto signals = fb.getSignals(); ASSERT_EQ(signals.getCount(), 1); const auto expectedComponentStatus = result ? "Warning" : "Error"; @@ -330,7 +330,7 @@ TEST_F(MqttSubscriberFbTest, RemovingNestedFunctionBlock) { auto config = PropertyObject(); config.addProperty(StringProperty(PROPERTY_NAME_SUB_TOPIC, String(buildTopicName()))); - ASSERT_NO_THROW(subFb = clientMqttFb.addFunctionBlock(SUB_FB_NAME, config)); + ASSERT_NO_THROW(subFb = mqttDevice.addFunctionBlock(SUB_FB_NAME, config)); ASSERT_EQ(subFb.getStatusContainer().getStatus("ComponentStatus"), Enumeration("ComponentStatusType", "Warning", daqInstance.getContext().getTypeManager())); } @@ -353,17 +353,17 @@ TEST_F(MqttSubscriberFbTest, TwoFbCreation) daq::FunctionBlockPtr fb; auto config = PropertyObject(); config.addProperty(StringProperty(PROPERTY_NAME_SUB_TOPIC, buildTopicName("0"))); - ASSERT_NO_THROW(fb = clientMqttFb.addFunctionBlock(SUB_FB_NAME, config)); + ASSERT_NO_THROW(fb = mqttDevice.addFunctionBlock(SUB_FB_NAME, config)); EXPECT_TRUE(waitStatusChange(1500, fb, Enumeration("ComponentStatusType", "Warning", daqInstance.getContext().getTypeManager()))); } { daq::FunctionBlockPtr fb; auto config = PropertyObject(); config.addProperty(StringProperty(PROPERTY_NAME_SUB_TOPIC, buildTopicName("1"))); - ASSERT_NO_THROW(fb = clientMqttFb.addFunctionBlock(SUB_FB_NAME, config)); + ASSERT_NO_THROW(fb = mqttDevice.addFunctionBlock(SUB_FB_NAME, config)); EXPECT_TRUE(waitStatusChange(1500, fb, Enumeration("ComponentStatusType", "Warning", daqInstance.getContext().getTypeManager()))); } - auto fbs = clientMqttFb.getFunctionBlocks(); + auto fbs = mqttDevice.getFunctionBlocks(); ASSERT_EQ(fbs.getCount(), 2u); } @@ -375,7 +375,7 @@ TEST_F(MqttSubscriberFbTest, PropertyChanged) auto config = PropertyObject(); auto topic = buildTopicName("0"); config.addProperty(StringProperty(PROPERTY_NAME_SUB_TOPIC, topic)); - ASSERT_NO_THROW(fb = clientMqttFb.addFunctionBlock(SUB_FB_NAME, config)); + ASSERT_NO_THROW(fb = mqttDevice.addFunctionBlock(SUB_FB_NAME, config)); EXPECT_TRUE(waitStatusChange(1500, fb, Enumeration("ComponentStatusType", "Warning", daqInstance.getContext().getTypeManager()))); auto subFb = reinterpret_cast(*fb); @@ -389,9 +389,9 @@ TEST_F(MqttSubscriberFbTest, JsonInit0) { StartUp(); daq::FunctionBlockPtr subFb; - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_SUB_JSON_CONFIG, String(VALID_JSON_1_TOPIC_0)); - ASSERT_NO_THROW(subFb = clientMqttFb.addFunctionBlock(SUB_FB_NAME, config)); + ASSERT_NO_THROW(subFb = mqttDevice.addFunctionBlock(SUB_FB_NAME, config)); ASSERT_EQ(subFb.getFunctionBlocks().getCount(), 3u); EXPECT_TRUE(waitStatusChange(1500, subFb, Enumeration("ComponentStatusType", "Warning", daqInstance.getContext().getTypeManager()))); auto lambda = [&](FunctionBlockPtr nestedFb, std::string value, std::string ts, std::string symbol) @@ -414,9 +414,9 @@ TEST_F(MqttSubscriberFbTest, JsonInit1) { StartUp(); daq::FunctionBlockPtr subFb; - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_SUB_JSON_CONFIG, String(VALID_JSON_1_TOPIC_1)); - ASSERT_NO_THROW(subFb = clientMqttFb.addFunctionBlock(SUB_FB_NAME, config)); + ASSERT_NO_THROW(subFb = mqttDevice.addFunctionBlock(SUB_FB_NAME, config)); ASSERT_EQ(subFb.getFunctionBlocks().getCount(), 3u); EXPECT_TRUE(waitStatusChange(1500, subFb, Enumeration("ComponentStatusType", "Warning", daqInstance.getContext().getTypeManager()))); auto lambda = [&](FunctionBlockPtr nestedFb, std::string value, std::string ts, std::string symbol) @@ -440,9 +440,9 @@ TEST_P(MqttSubscriberFbConfigPTest, JsonWrongInit) const auto configJson = GetParam(); StartUp(); daq::FunctionBlockPtr subFb; - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_SUB_JSON_CONFIG, String(configJson)); - ASSERT_NO_THROW(subFb = clientMqttFb.addFunctionBlock(SUB_FB_NAME, config)); + ASSERT_NO_THROW(subFb = mqttDevice.addFunctionBlock(SUB_FB_NAME, config)); EXPECT_EQ(subFb.getFunctionBlocks().getCount(), 0u); ASSERT_TRUE(waitStatusChange(1500, subFb, Enumeration("ComponentStatusType", "Error", daqInstance.getContext().getTypeManager()))); ASSERT_TRUE(waitStatusTheSame(1500, subFb, Enumeration("ComponentStatusType", "Error", daqInstance.getContext().getTypeManager()))); @@ -464,9 +464,9 @@ TEST_P(MqttSubscriberFbConfigFilePTest, JsonInitFromFile) const auto configJson = GetParam(); StartUp(); daq::FunctionBlockPtr subFb; - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_SUB_JSON_CONFIG_FILE, String(configJson)); - ASSERT_NO_THROW(subFb = clientMqttFb.addFunctionBlock(SUB_FB_NAME, config)); + ASSERT_NO_THROW(subFb = mqttDevice.addFunctionBlock(SUB_FB_NAME, config)); EXPECT_TRUE(waitStatusChange(1500, subFb, Enumeration("ComponentStatusType", "Warning", daqInstance.getContext().getTypeManager()))); } @@ -481,9 +481,9 @@ TEST_F(MqttSubscriberFbTest, JsonInitFromFileWithChecking) { StartUp(); daq::FunctionBlockPtr subFb; - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_SUB_JSON_CONFIG_FILE, String("data/public-example0.json")); - ASSERT_NO_THROW(subFb = clientMqttFb.addFunctionBlock(SUB_FB_NAME, config)); + ASSERT_NO_THROW(subFb = mqttDevice.addFunctionBlock(SUB_FB_NAME, config)); EXPECT_TRUE(waitStatusChange(1500, subFb, Enumeration("ComponentStatusType", "Warning", daqInstance.getContext().getTypeManager()))); ASSERT_EQ(subFb.getFunctionBlocks().getCount(), 3u); auto lambda = [&](FunctionBlockPtr nestedFb, std::string value, std::string ts, std::string symbol) @@ -573,9 +573,9 @@ TEST_F(MqttSubscriberFbTest, JsonInitFromFileWrongPath) { StartUp(); daq::FunctionBlockPtr subFb; - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_SUB_JSON_CONFIG_FILE, String("/justWrongPath/wrongFile.txt")); - ASSERT_NO_THROW(subFb = clientMqttFb.addFunctionBlock(SUB_FB_NAME, config)); + ASSERT_NO_THROW(subFb = mqttDevice.addFunctionBlock(SUB_FB_NAME, config)); EXPECT_EQ(subFb.getFunctionBlocks().getCount(), 0u); ASSERT_TRUE(waitStatusChange(1500, subFb, Enumeration("ComponentStatusType", "Error", daqInstance.getContext().getTypeManager()))); ASSERT_TRUE(waitStatusTheSame(1500, subFb, Enumeration("ComponentStatusType", "Error", daqInstance.getContext().getTypeManager()))); @@ -652,12 +652,12 @@ TEST_F(MqttSubscriberFbTest, WaitingData) StartUp(); - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_SUB_TOPIC, topic); config.setPropertyValue(PROPERTY_NAME_SUB_PREVIEW_SIGNAL, True); config.setPropertyValue(PROPERTY_NAME_SUB_DATA_TIMEOUT, 200); - auto rawFB = clientMqttFb.addFunctionBlock(SUB_FB_NAME, config); + auto rawFB = mqttDevice.addFunctionBlock(SUB_FB_NAME, config); ASSERT_TRUE(waitStatusChange(1000, rawFB, Enumeration("ComponentStatusType", "Warning", daqInstance.getContext().getTypeManager()))); EXPECT_NE(rawFB.getStatusContainer().getStatusMessage("ComponentStatus").toStdString().find("Waiting for data"), std::string::npos); @@ -702,11 +702,11 @@ TEST_F(MqttSubscriberFbTest, CheckRawFbFullDataTransfer) StartUp(); - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_SUB_TOPIC, topic); config.setPropertyValue(PROPERTY_NAME_SUB_PREVIEW_SIGNAL, True); - auto singal = clientMqttFb.addFunctionBlock(SUB_FB_NAME, config).getSignals()[0]; + auto singal = mqttDevice.addFunctionBlock(SUB_FB_NAME, config).getSignals()[0]; auto reader = daq::PacketReader(singal); MqttAsyncClientWrapper publisher("testPublisherId"); @@ -754,10 +754,10 @@ TEST_F(MqttSubscriberFbTest, CheckRawFbFullDataTransferWithReconfiguring) StartUp(); - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_SUB_TOPIC, topic0); config.setPropertyValue(PROPERTY_NAME_SUB_PREVIEW_SIGNAL, True); - auto rawFB = clientMqttFb.addFunctionBlock(SUB_FB_NAME, config); + auto rawFB = mqttDevice.addFunctionBlock(SUB_FB_NAME, config); auto singal = rawFB.getSignals()[0]; auto reader = daq::PacketReader(singal); @@ -820,12 +820,12 @@ TEST_F(MqttSubscriberFbTest, DomainDataPacketWithTheSameTS) StartUp(); - auto config = clientMqttFb.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); + auto config = mqttDevice.getAvailableFunctionBlockTypes().get(SUB_FB_NAME).createDefaultConfig(); config.setPropertyValue(PROPERTY_NAME_SUB_TOPIC, topic); config.setPropertyValue(PROPERTY_NAME_SUB_PREVIEW_SIGNAL, True); config.setPropertyValue(PROPERTY_NAME_SUB_PREVIEW_SIGNAL_TS_MODE, static_cast(SDSM::SystemTime)); config.setPropertyValue(PROPERTY_NAME_SUB_DATA_TIMEOUT, 0); - auto fb = clientMqttFb.addFunctionBlock(SUB_FB_NAME, config); + auto fb = mqttDevice.addFunctionBlock(SUB_FB_NAME, config); auto getTime = []() { return duration_cast(system_clock::now().time_since_epoch()).count(); }; auto getStatusMsg = [&]() { return fb.getStatusContainer().getStatusMessage("ComponentStatus").toStdString(); };