Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 10 additions & 2 deletions CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ option(${REPO_OPTION_PREFIX}_ENABLE_EXAMPLE_APP "Enable ${REPO_NAME} example app
option(${REPO_OPTION_PREFIX}_ENABLE_TESTS "Enable ${REPO_NAME} testing" ${PROJECT_IS_TOP_LEVEL})
option(${REPO_OPTION_PREFIX}_ENABLE_CLIENT "Enable ${REPO_NAME} client module" ${PROJECT_IS_TOP_LEVEL})
option(${REPO_OPTION_PREFIX}_ENABLE_SERVER "Enable ${REPO_NAME} server module" ${PROJECT_IS_TOP_LEVEL})
option(${REPO_OPTION_PREFIX}_ENABLE_TLS "Enable ${REPO_NAME} TLS (daq.lts://) streaming channel" ON)

opendaq_common_compile_targets_settings()
opendaq_setup_compiler_flags(${REPO_OPTION_PREFIX})
Expand All @@ -40,6 +41,12 @@ if (CMAKE_CXX_COMPILER_ID MATCHES "Clang|AppleClang" AND CMAKE_CXX_COMPILER_VERS
add_compile_options(-Wno-unknown-warning-option)
endif()

if (${REPO_OPTION_PREFIX}_ENABLE_TLS)
message(STATUS "TLS streaming channel in ${REPO_NAME} is ENABLED")
else()
message(STATUS "TLS streaming channel in ${REPO_NAME} is DISABLED (no OpenSSL dependency)")
endif()

if (${REPO_OPTION_PREFIX}_ENABLE_TESTS)
message(STATUS "Unit tests in ${REPO_NAME} are ENABLED")
enable_testing()
Expand Down Expand Up @@ -77,8 +84,9 @@ add_subdirectory(external)
add_subdirectory(shared)
add_subdirectory(modules)

# End-to-end integration tests need both modules present
if (${REPO_OPTION_PREFIX}_ENABLE_TESTS AND ${REPO_OPTION_PREFIX}_ENABLE_CLIENT AND ${REPO_OPTION_PREFIX}_ENABLE_SERVER)
# End-to-end integration tests need both modules present, and only cover the TLS channel
if (${REPO_OPTION_PREFIX}_ENABLE_TESTS AND ${REPO_OPTION_PREFIX}_ENABLE_CLIENT
AND ${REPO_OPTION_PREFIX}_ENABLE_SERVER AND ${REPO_OPTION_PREFIX}_ENABLE_TLS)
add_subdirectory(tests)
endif()

4 changes: 3 additions & 1 deletion external/ws-streaming/CMakeLists.txt
Original file line number Diff line number Diff line change
@@ -1,10 +1,12 @@
set(WS_STREAMING_INSTALL OFF)

set(WS_STREAMING_ENABLE_TLS ${${REPO_OPTION_PREFIX}_ENABLE_TLS})

opendaq_dependency(
NAME ws-streaming
REQUIRED_VERSION 3.2.0
GIT_REPOSITORY https://github.com/openDAQ/ws-streaming
GIT_REF main
GIT_REF tls-optional
GIT_SHALLOW ON
EXPECT_TARGET ws-streaming::ws-streaming
)
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,14 @@ class WebsocketStreamingClientModule final : public Module
std::string path;
};

// Pick the plain or the secure variant according to the connection string. The daq.lts://
// counterparts only exist in a build with the TLS channel, so these keep the preprocessor
// out of onCreateDevice() and onCreateStreaming().
static PropertyObjectPtr createDefaultDeviceConfig(const StringPtr& connectionString);
static PropertyObjectPtr createDefaultStreamingConfig(const StringPtr& connectionString);
static DeviceTypePtr createDeviceType(const StringPtr& connectionString);
static StreamingTypePtr createStreamingType(const StringPtr& connectionString);

DAQ_WS_STREAM_CL_MODULE_API static StringPtr createUrlConnectionString(bool secureType,
const StringPtr& host,
const IntegerPtr& port,
Expand All @@ -60,6 +68,7 @@ class WebsocketStreamingClientModule final : public Module
DAQ_WS_STREAM_CL_MODULE_API static StringPtr formNewStyleConnectionString(const StringPtr& connectionString);
DAQ_WS_STREAM_CL_MODULE_API static DeviceInfoPtr populateDiscoveredDevice(const discovery::MdnsDiscoveredDevice& discoveredDevice);
DAQ_WS_STREAM_CL_MODULE_API static bool isSecureConnection(const std::string& connectionString);
DAQ_WS_STREAM_CL_MODULE_API static bool isSupportedServiceName(const std::string& serviceName);

std::mutex sync;
size_t deviceIndex;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -54,15 +54,27 @@ WebsocketStreamingClientModule::WebsocketStreamingClientModule(ContextPtr contex
, deviceIndex(0)
, discoveryClient({CONST_SERVICE_CAPABILITY})
{
#if DAQMODULES_LT_STREAMING_ENABLE_TLS
discoveryClient.initMdnsClient(List<IString>(CONST_LT_SERVICE_NAME, CONST_LTS_SERVICE_NAME, CONST_WS_SERVICE_NAME));
#else
discoveryClient.initMdnsClient(List<IString>(CONST_LT_SERVICE_NAME, CONST_WS_SERVICE_NAME));
#endif
loggerComponent = this->context.getLogger().getOrAddComponent("StreamingLTClient");
}

ListPtr<IDeviceInfo> WebsocketStreamingClientModule::onGetAvailableDevices()
{
auto availableDevices = List<IDeviceInfo>();
for (const auto& device : discoveryClient.discoverMdnsDevices())
{
if (!isSupportedServiceName(device.serviceName))
{
LOG_D("Ignoring discovered service \"{}\": not supported by this module", device.serviceName)
continue;
}

availableDevices.pushBack(populateDiscoveredDevice(device));
}
return availableDevices;
}

Expand All @@ -72,11 +84,14 @@ DictPtr<IString, IDeviceType> WebsocketStreamingClientModule::onGetAvailableDevi

const auto websocketDeviceType = WsStreamingDevice::createNewType();
const auto oldWebsocketDeviceType = WsStreamingDevice::createOldType();
const auto secureWebsocketDeviceType = WsStreamingDevice::createNewSecureType();

result.set(websocketDeviceType.getId(), websocketDeviceType);
result.set(oldWebsocketDeviceType.getId(), oldWebsocketDeviceType);

#if DAQMODULES_LT_STREAMING_ENABLE_TLS
const auto secureWebsocketDeviceType = WsStreamingDevice::createNewSecureType();
result.set(secureWebsocketDeviceType.getId(), secureWebsocketDeviceType);
#endif

return result;
}
Expand All @@ -86,10 +101,12 @@ DictPtr<IString, IStreamingType> WebsocketStreamingClientModule::onGetAvailableS
auto result = Dict<IString, IStreamingType>();

auto websocketStreamingType = WsStreaming::createType();
auto secureWebsocketStreamingType = WsStreaming::createSecureType();

result.set(websocketStreamingType.getId(), websocketStreamingType);

#if DAQMODULES_LT_STREAMING_ENABLE_TLS
auto secureWebsocketStreamingType = WsStreaming::createSecureType();
result.set(secureWebsocketStreamingType.getId(), secureWebsocketStreamingType);
#endif

return result;
}
Expand All @@ -115,20 +132,17 @@ DevicePtr WebsocketStreamingClientModule::onCreateDevice(const StringPtr& connec

PropertyObjectPtr deviceConfig = config;
if (!deviceConfig.assigned())
deviceConfig = isSecureConnection(formedConnectionStr) ? WsStreamingDevice::createDefaultSecureConfig()
: WsStreamingDevice::createDefaultConfig();
deviceConfig = createDefaultDeviceConfig(formedConnectionStr);

std::scoped_lock lock(sync);

std::string localId = fmt::format("websocket_pseudo_device{}", deviceIndex++);
auto deviceType =
isSecureConnection(formedConnectionStr) ? WsStreamingDevice::createNewSecureType() : WsStreamingDevice::createNewType();
auto deviceType = createDeviceType(formedConnectionStr);
checkErrorInfo(deviceType.asPtr<IComponentTypePrivate>()->setModuleInfo(moduleInfo));
auto device = createWithImplementation<IDevice, WsStreamingDevice>(context, parent, localId, formedConnectionStr, deviceType, deviceConfig);

// Set the connection info for the device
const auto wsStreamingType =
(isSecureConnection(formedConnectionStr)) ? WsStreaming::createSecureType() : WsStreaming::createType();
const auto wsStreamingType = createStreamingType(formedConnectionStr);

ServerCapabilityConfigPtr connectionInfo = device.getInfo().getConfigurationConnectionInfo();
connectionInfo.setProtocolId(wsStreamingType.getId());
Expand Down Expand Up @@ -174,9 +188,7 @@ StreamingPtr WebsocketStreamingClientModule::onCreateStreaming(const StringPtr&

PropertyObjectPtr streamingConfig = config;
if (!streamingConfig.assigned())
streamingConfig = isSecureConnection(formNewStyleConnectionString(connectionString).toStdString())
? WsStreaming::createDefaultSecureConfig()
: WsStreaming::createDefaultConfig();
streamingConfig = createDefaultStreamingConfig(formNewStyleConnectionString(connectionString));

const StringPtr str = formConnectionString(connectionString, streamingConfig);
return createWithImplementation<IStreaming, WsStreaming>(str, context, streamingConfig);
Expand All @@ -185,8 +197,13 @@ StreamingPtr WebsocketStreamingClientModule::onCreateStreaming(const StringPtr&
Bool WebsocketStreamingClientModule::onCompleteServerCapability(const ServerCapabilityPtr& source, const ServerCapabilityConfigPtr& target)
{
const auto protoId = target.getProtocolId();
#if DAQMODULES_LT_STREAMING_ENABLE_TLS
if (protoId != CONST_LT_STREAMING_ID && protoId != CONST_LTS_STREAMING_ID)
return false;
#else
if (protoId != CONST_LT_STREAMING_ID)
return false;
#endif

if (source.getConnectionType() != "TCP/IP")
return false;
Expand Down Expand Up @@ -243,6 +260,46 @@ Bool WebsocketStreamingClientModule::onCompleteServerCapability(const ServerCapa
return true;
}

PropertyObjectPtr WebsocketStreamingClientModule::createDefaultDeviceConfig(const StringPtr& connectionString)
{
#if DAQMODULES_LT_STREAMING_ENABLE_TLS
if (isSecureConnection(connectionString))
return WsStreamingDevice::createDefaultSecureConfig();
#endif

return WsStreamingDevice::createDefaultConfig();
}

PropertyObjectPtr WebsocketStreamingClientModule::createDefaultStreamingConfig(const StringPtr& connectionString)
{
#if DAQMODULES_LT_STREAMING_ENABLE_TLS
if (isSecureConnection(connectionString))
return WsStreaming::createDefaultSecureConfig();
#endif

return WsStreaming::createDefaultConfig();
}

DeviceTypePtr WebsocketStreamingClientModule::createDeviceType(const StringPtr& connectionString)
{
#if DAQMODULES_LT_STREAMING_ENABLE_TLS
if (isSecureConnection(connectionString))
return WsStreamingDevice::createNewSecureType();
#endif

return WsStreamingDevice::createNewType();
}

StreamingTypePtr WebsocketStreamingClientModule::createStreamingType(const StringPtr& connectionString)
{
#if DAQMODULES_LT_STREAMING_ENABLE_TLS
if (isSecureConnection(connectionString))
return WsStreaming::createSecureType();
#endif

return WsStreaming::createType();
}

StringPtr WebsocketStreamingClientModule::createUrlConnectionString(bool secureType,
const StringPtr& host,
const IntegerPtr& port,
Expand Down Expand Up @@ -330,12 +387,14 @@ DeviceInfoPtr WebsocketStreamingClientModule::populateDiscoveredDevice(const Mdn
streamingType = WsStreaming::createType();
deviceType = WsStreamingDevice::createNewType();
}
#if DAQMODULES_LT_STREAMING_ENABLE_TLS
else if (discoveredDevice.serviceName == CONST_LTS_SERVICE_NAME)
{
isSecure = true;
streamingType = WsStreaming::createSecureType();
deviceType = WsStreamingDevice::createNewSecureType();
}
#endif
else
{
DAQ_THROW_EXCEPTION(InvalidParameterException,
Expand Down Expand Up @@ -387,4 +446,17 @@ bool WebsocketStreamingClientModule::isSecureConnection(const std::string& conne
return connectionString.find(securePrefix) != std::string::npos;
}

bool WebsocketStreamingClientModule::isSupportedServiceName(const std::string& serviceName)
{
if (serviceName == CONST_LT_SERVICE_NAME || serviceName == CONST_WS_SERVICE_NAME)
return true;

#if DAQMODULES_LT_STREAMING_ENABLE_TLS
if (serviceName == CONST_LTS_SERVICE_NAME)
return true;
#endif

return false;
}

END_NAMESPACE_OPENDAQ_WEBSOCKET_STREAMING_CLIENT_MODULE
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,10 @@ set(TEST_SOURCES test_websocket_streaming_client_module.cpp
test_app.cpp
)

if (${REPO_OPTION_PREFIX}_ENABLE_TLS)
list(APPEND TEST_SOURCES test_websocket_streaming_client_module_tls.cpp)
endif()

add_executable(${TEST_APP} ${TEST_SOURCES}
)

Expand Down
Loading
Loading