diff --git a/.github/workflows/validate.yml b/.github/workflows/validate.yml index ff3ad657..f162a1aa 100644 --- a/.github/workflows/validate.yml +++ b/.github/workflows/validate.yml @@ -19,6 +19,10 @@ concurrency: jobs: install-lint-test: strategy: + # Let each platform report its true status: one platform's failure + # must not cancel the others (the macOS legs are expected to fail + # until the Phase 1.2 compiler upgrade lands; Linux is the gate). + fail-fast: false matrix: os: [macos-14, macos-13, ubuntu-24.04, linux-arm64-runner] runs-on: ${{ matrix.os }} diff --git a/.gitignore b/.gitignore index 357f1609..4af7fe19 100644 --- a/.gitignore +++ b/.gitignore @@ -1,5 +1,6 @@ vcpkg_installed *.cache* +*~ .vscode .trunk setenvs.sh diff --git a/packages/streamr-json/test/unit/toJsonTest.cpp~ b/packages/streamr-json/test/unit/toJsonTest.cpp~ deleted file mode 100644 index 836572d4..00000000 --- a/packages/streamr-json/test/unit/toJsonTest.cpp~ +++ /dev/null @@ -1,367 +0,0 @@ -#include -#include -#include -#include -#include -#include -#include - -#include - -#include "TestClass.hpp" -#include "WeatherData.hpp" -#include "streamr-json/toJson.hpp" - -using streamr::json::toJson; - -// NOLINTBEGIN(readability-magic-numbers) - -class ToJsonTest : public ::testing::Test { -protected: - WeatherData weatherData; // NOLINT - TestClass testClass; // NOLINT - json weatherJson; // NOLINT - - void SetUp() override { - weatherData = { - .dataId = 0, - .dataLabel = "Test data", - .dataByCountry = { - {"Finland", - { - {.locality = "Helsinki", - .temperatures = - { - {.temperature = 23.4, .timestamp = 1}, // NOLINT - {.temperature = 24.5, .timestamp = 1000} // NOLINT - }}, - }}, - {"Sweden", - {{.locality = "Stockholm", - .temperatures = { - {.temperature = 22.0, .timestamp = 1}, // NOLINT - {.temperature = 21.6, .timestamp = 1000} // NOLINT - }}}}}}; - weatherJson = R"({ - "dataId": 0, - "dataLabel": "Test data", - "dataByCountry": { - "Finland": [ - { - "locality": "Helsinki", - "temperatures": [ - {"temperature": 23.4, "timestamp": 1}, - {"temperature": 24.5, "timestamp": 1000} - ] - } - ], - "Sweden": [ - { - "locality": "Stockholm", - "temperatures": [ - {"temperature": 22.0, "timestamp": 1}, - {"temperature": 21.6, "timestamp": 1000} - ] - } - ] - } - })"_json; - } -}; - -TEST_F(ToJsonTest, TestWeatherDataToJson) { - EXPECT_EQ(toJson(weatherData), weatherJson); -} - -TEST_F(ToJsonTest, TestInitializerListToJson) { - EXPECT_EQ( - toJson( - {{"pi", 3.141}, - {"happy", true}, - {"name", "Niels"}, - {"nothing", nullptr}, - {"answer", {{"everything", 42}}}, - {"list", {1, 0, 2}}, - {"object", {{"currency", "USD"}, {"value", 42.99}}}}), - json( - {{"pi", 3.141}, - {"happy", true}, - {"name", "Niels"}, - {"nothing", nullptr}, - {"answer", {{"everything", 42}}}, - {"list", {1, 0, 2}}, - {"object", {{"currency", "USD"}, {"value", 42.99}}}})); -} - -TEST_F(ToJsonTest, TestInitializerListWithStructContentToJson) { - EXPECT_EQ(toJson({{"data", weatherData}}), json({{"data", weatherJson}})); -} - -TEST_F(ToJsonTest, TestInitializerListWithStructContentEmptyInputToJson) { - json expectedJson = { - {"data", {{"dataId", 0}, {"dataLabel", ""}, {"dataByCountry", {}}}}}; - WeatherData emptyData{}; - EXPECT_EQ(toJson({{"data", emptyData}}), expectedJson); -} - -TEST_F(ToJsonTest, TestInitializerListWithStructContentNullInputToJson) { - json expectedJson = {{"data", nullptr}}; - WeatherData* nullData = nullptr; - EXPECT_EQ(toJson({{"data", nullptr}}), expectedJson); -} - -TEST_F(ToJsonTest, TestInitializerListWithStructContentZeroInputToJson) { - json expectedJson = { - {"data", {{"dataId", 0}, {"dataLabel", ""}, {"dataByCountry", {}}}}}; - WeatherData zeroData; - zeroData.dataId = 0; - zeroData.dataLabel = ""; - zeroData.dataByCountry = {}; - EXPECT_EQ(toJson({{"data", zeroData}}), expectedJson); -}; - -TEST_F(ToJsonTest, TestInitializerListWithStructContent) { - EXPECT_EQ(toJson({{"data", weatherData}}), json({{"data", weatherJson}})); -} - -TEST_F(ToJsonTest, TestIntToJson) { - int a = 10; - EXPECT_EQ(toJson(a), 10); -} - -TEST_F(ToJsonTest, TestDoubleToJson) { - double b = 20.5; - EXPECT_EQ(toJson(b), 20.5); -} - -TEST_F(ToJsonTest, TestCharToJson) { - char c = 'c'; - EXPECT_EQ( - toJson(c), 99); // char is a number type in C++, 'c' is 99 in ASCII -} - -TEST_F(ToJsonTest, TestStringToJson) { - std::string d = "Hello, World!"; - EXPECT_EQ(toJson(d), "Hello, World!"); -} - -TEST_F(ToJsonTest, TestBoolToJson) { - bool e = true; - EXPECT_EQ(toJson(e), true); -} - -TEST_F(ToJsonTest, TestVectorToJson) { - std::vector f = {1, 2, 3, 4, 5}; // NOLINT - EXPECT_EQ(toJson(f), json({1, 2, 3, 4, 5})); -} - -TEST_F(ToJsonTest, TestMapToJson) { - std::map g = {{"one", 1}, {"two", 2}, {"three", 3}}; - // std::map orders keys alphabetically - EXPECT_EQ(toJson(g), json({{"one", 1}, {"three", 3}, {"two", 2}})); -} - -TEST_F(ToJsonTest, TestPairToJson) { - std::pair h = {1, "one"}; - EXPECT_EQ(toJson(h), json({1, "one"})); -} - -TEST_F(ToJsonTest, TestSetToJson) { - std::set i = {1, 2, 3, 4, 5}; // NOLINT - EXPECT_EQ(toJson(i), json({1, 2, 3, 4, 5})); -} - -TEST_F(ToJsonTest, TestArrayToJson) { - std::array j = {1, 2, 3, 4, 5}; // NOLINT - EXPECT_EQ(toJson(j), json({1, 2, 3, 4, 5})); -} - -TEST_F(ToJsonTest, TestFloatToJson) { - float k = 30.5f; // NOLINT - EXPECT_EQ(toJson(k), 30.5); -} - -TEST_F(ToJsonTest, TestLongToJson) { - long l = 1000000L; // NOLINT - EXPECT_EQ(toJson(l), 1000000); -} - -TEST_F(ToJsonTest, TestShortToJson) { - short m = 10; // NOLINT - EXPECT_EQ(toJson(m), 10); -} - -TEST_F(ToJsonTest, TestUnsignedIntToJson) { - unsigned int n = 20; // NOLINT - EXPECT_EQ(toJson(n), 20); -} - -TEST_F(ToJsonTest, TestLongDoubleToJson) { - long double o = 30.5L; // NOLINT - EXPECT_EQ(toJson(o), 30.5); -} - -TEST_F(ToJsonTest, TestUnsignedLongToJson) { - unsigned long p = 1000000UL; // NOLINT - EXPECT_EQ(toJson(p), 1000000); -} - -TEST_F(ToJsonTest, TestUnsignedShortToJson) { - unsigned short q = 10; // NOLINT - EXPECT_EQ(toJson(q), 10); -} - -TEST_F(ToJsonTest, TestListToJson) { - std::list r = {1, 2, 3, 4, 5}; // NOLINT - EXPECT_EQ(toJson(r), json({1, 2, 3, 4, 5})); -} - -TEST_F(ToJsonTest, TestDequeToJson) { - std::deque s = {1, 2, 3, 4, 5}; // NOLINT - EXPECT_EQ(toJson(s), json({1, 2, 3, 4, 5})); -} - -TEST_F(ToJsonTest, TestForwardListToJson) { - std::forward_list t = {1, 2, 3, 4, 5}; // NOLINT - EXPECT_EQ(toJson(t), json({1, 2, 3, 4, 5})); -} - -TEST_F(ToJsonTest, TestClassToJson) { - json expectedJson = {{"a", testClass.getA()}, {"b", testClass.getB()}}; - EXPECT_EQ(toJson(testClass), expectedJson); -} - -TEST_F(ToJsonTest, TestEmptyStructToJson) { - struct EmptyStruct { - } emptyStruct; - EXPECT_EQ(toJson(emptyStruct), json({})); -} - -TEST_F(ToJsonTest, TestSpecialCharactersToJson) { - std::string specialCharacters = "!@#$%^&*()"; - EXPECT_EQ(toJson(specialCharacters), "!@#$%^&*()"); -} - -TEST_F(ToJsonTest, TestNullPointerToJson) { - int* nullPointer = nullptr; - json j(nullptr); - EXPECT_EQ(toJson(nullPointer), j); -} - -TEST_F(ToJsonTest, TestNonNullPointerToJson) { - int a = 1; - int* nonNullPointer = &a; - - json j(*nonNullPointer); - - EXPECT_EQ(toJson(nonNullPointer), j); -} - -TEST_F(ToJsonTest, TestWeatherDataSmartPointersToJson) { - WeatherDataSmartPointers weatherDataSmartPointers; - weatherDataSmartPointers.dataId = 0; - weatherDataSmartPointers.dataLabel = - std::make_shared("Test data"); - - auto finlandDataSample = std::make_shared(); - finlandDataSample->locality = "Helsinki"; - finlandDataSample->temperatures.push_back({23.4, 1}); - finlandDataSample->temperatures.push_back({24.5, 1000}); - weatherDataSmartPointers.dataByCountry["Finland"].push_back( - finlandDataSample); - - auto swedenDataSample = std::make_shared(); - swedenDataSample->locality = "Stockholm"; - swedenDataSample->temperatures.push_back({22.0, 1}); - swedenDataSample->temperatures.push_back({21.6, 1000}); - weatherDataSmartPointers.dataByCountry["Sweden"].push_back( - swedenDataSample); - - EXPECT_EQ(toJson(weatherDataSmartPointers), weatherJson); -} - -TEST_F(ToJsonTest, TestWeatherDataRegularPointersToJson) { - WeatherDataRegularPointers weatherDataRegularPointers; - weatherDataRegularPointers.dataId = 0; - weatherDataRegularPointers.dataLabel = new std::string("Test data"); - auto* finlandSample = new DataSample; - finlandSample->locality = "Helsinki"; - finlandSample->temperatures.push_back({23.4, 1}); - finlandSample->temperatures.push_back({24.5, 1000}); - weatherDataRegularPointers.dataByCountry["Finland"].push_back( - finlandSample); - - auto* swedenSample = new DataSample; - swedenSample->locality = "Stockholm"; - swedenSample->temperatures.push_back({22.0, 1}); - swedenSample->temperatures.push_back({21.6, 1000}); - weatherDataRegularPointers.dataByCountry["Sweden"].push_back(swedenSample); - json expectedJson = R"( - { - "dataId": 0, - "dataLabel": "Test data", - "dataByCountry": { - "Finland": [ - { - "locality": "Helsinki", - "temperatures": [ - {"temperature": 23.4, "timestamp": 1}, - {"temperature": 24.5, "timestamp": 1000} - ] - } - ], - "Sweden": [ - { - "locality": "Stockholm", - "temperatures": [ - {"temperature": 22.0, "timestamp": 1}, - {"temperature": 21.6, "timestamp": 1000} - ] - } - ] - } - } - )"_json; - /* - json expectedJson = { - {"dataId", 0}, - {"dataLabel", "Test data"}, - {"dataByCountry", - {"Finland", { - { - { - {"locality", "Helsinki"}, - {"temperatures", - { - {{"temperature", 23.4}, {"timestamp", - std::time(nullptr)}}, - {{"temperature", 24.5}, {"timestamp", - std::time(nullptr) + 1000}} - } - } - } - }}}}}}, - {"Sweden", { - { - { - {"locality", "Stockholm"}, - {"temperatures", - { - {{"temperature", 22.0}, {"timestamp", std::time(nullptr)}}, - {{"temperature", 21.6}, {"timestamp", std::time(nullptr) + - 1000}} - }}}}}} - };*/ - - EXPECT_EQ(toJson(weatherDataRegularPointers), expectedJson); - - delete weatherDataRegularPointers.dataLabel; - for (auto& sample : weatherDataRegularPointers.dataByCountry["Finland"]) { - delete sample; - } - for (auto& sample : weatherDataRegularPointers.dataByCountry["Sweden"]) { - delete sample; - } -} - -// NOLINTEND(readability-magic-numbers) diff --git a/packages/streamr-libstreamrproxyclient/CMakeLists.txt~ b/packages/streamr-libstreamrproxyclient/CMakeLists.txt~ deleted file mode 100644 index 22fd0222..00000000 --- a/packages/streamr-libstreamrproxyclient/CMakeLists.txt~ +++ /dev/null @@ -1,93 +0,0 @@ -cmake_minimum_required(VERSION 3.22) -set(CMAKE_POLICY_DEFAULT_CMP0077 NEW) - -include(homebrewClang.cmake) -set(CMAKE_CXX_STANDARD 26) - -set(CMAKE_CXX_USE_RESPONSE_FILE_FOR_INCLUDES Off) -set(CMAKE_EXPORT_COMPILE_COMMANDS ON) - -message(STATUS "CMAKE_CURRENT_SOURCE_DIR: ${CMAKE_CURRENT_SOURCE_DIR}") - -# make the current package a monorepo package -include(${CMAKE_CURRENT_SOURCE_DIR}/monorepoPackage.cmake) - -set(CMAKE_TOOLCHAIN_FILE "$ENV{VCPKG_ROOT}/scripts/buildsystems/vcpkg.cmake") - -project(streamr-streamrproxyclient CXX) -add_library(streamrproxyclient SHARED - src/streamrproxyclient.cpp - include/streamrproxyclient.h -) - -find_package(streamr-trackerless-network CONFIG REQUIRED) -find_package(streamr-dht CONFIG REQUIRED) -find_package(streamr-logger CONFIG REQUIRED) -find_package(streamr-utils CONFIG REQUIRED) -find_package(ada CONFIG REQUIRED) -find_package(cpptrace CONFIG REQUIRED) - -target_include_directories(streamrproxyclient PUBLIC include) -target_link_libraries(streamrproxyclient - PRIVATE streamr::streamr-trackerless-network - PRIVATE streamr::streamr-logger - PRIVATE streamr::streamr-dht - PRIVATE streamr::streamr-utils - PRIVATE ada::ada - PRIVATE Folly::folly -) - -if (IOS) - message(STATUS "Building dylib for iOS") - target_link_libraries(streamrproxyclient ${FOUNDATION_LIBRARY}) - target_compile_definitions(streamrproxyclient PUBLIC IS_BUILDING_SHARED) - set_target_properties(streamrproxyclient PROPERTIES - FRAMEWORK TRUE - FRAMEWORK_VERSION A - MACOSX_FRAMEWORK_IDENTIFIER network.streamr.streamrproxyclient - VERSION 1.0.0 - SOVERSION 1.0.0 - PUBLIC_HEADER include/streamrproxyclient.h - XCODE_ATTRIBUTE_CODE_SIGN_IDENTITY "Petri Savolainen" - ) - install(CODE "execute_process(COMMAND ${CMAKE_CURRENT_LIST_DIR}/createiosframework.sh)") -endif() - -if (NOT IOS) - enable_testing() - find_package(GTest CONFIG REQUIRED) - find_package(folly CONFIG REQUIRED) - add_executable(streamr-streamrproxyclient-test-integration test/integration/StreamrProxyClientTest.cpp) - target_link_directories(streamr-streamrproxyclient-test-integration PUBLIC ${CMAKE_CURRENT_LIST_DIR}/build) - target_include_directories(streamr-streamrproxyclient-test-integration PUBLIC ${CMAKE_CURRENT_LIST_DIR}/include) - target_link_libraries(streamr-streamrproxyclient-test-integration - PUBLIC streamrproxyclient - PRIVATE streamr::streamr-logger - PUBLIC GTest::gtest - PUBLIC GTest::gtest_main - PUBLIC cpptrace::cpptrace - ) - - add_executable(streamr-streamrproxyclient-test-end-to-end test/end-to-end/PublishToTsServerTest.cpp) - target_link_directories(streamr-streamrproxyclient-test-end-to-end PUBLIC ${CMAKE_CURRENT_LIST_DIR}/build) - target_include_directories(streamr-streamrproxyclient-test-end-to-end PUBLIC ${CMAKE_CURRENT_LIST_DIR}/include) - target_link_libraries(streamr-streamrproxyclient-test-end-to-end - PUBLIC streamrproxyclient - PRIVATE streamr::streamr-logger - PUBLIC GTest::gtest - PUBLIC GTest::gtest_main - PUBLIC cpptrace::cpptrace - ) - - if (NOT (${VCPKG_TARGET_TRIPLET} MATCHES "android")) - include(GoogleTest) - gtest_discover_tests(streamr-streamrproxyclient-test-integration) - - add_test( - NAME ts-end-to-end-test - COMMAND "${CMAKE_CURRENT_LIST_DIR}/run-ts-end-to-end-tests.sh" - WORKING_DIRECTORY "${CMAKE_CURRENT_LIST_DIR}" - ) - endif() -endif() - diff --git a/packages/streamr-libstreamrproxyclient/wrappers/go b/packages/streamr-libstreamrproxyclient/wrappers/go index cdb68c95..55f17ba3 160000 --- a/packages/streamr-libstreamrproxyclient/wrappers/go +++ b/packages/streamr-libstreamrproxyclient/wrappers/go @@ -1 +1 @@ -Subproject commit cdb68c95644dd2642a046f576cf777f490c2cce1 +Subproject commit 55f17ba3850938aa7539bb0f22409b4abe1b15a4 diff --git a/packages/streamr-proto-rpc/CMakeLists.txt~ b/packages/streamr-proto-rpc/CMakeLists.txt~ deleted file mode 100644 index c3ca7d63..00000000 --- a/packages/streamr-proto-rpc/CMakeLists.txt~ +++ /dev/null @@ -1,170 +0,0 @@ -cmake_minimum_required(VERSION 3.22) -set(CMAKE_POLICY_DEFAULT_CMP0077 NEW) - -include(homebrewClang.cmake) - -set(CMAKE_CXX_STANDARD 26) -set(CMAKE_EXPORT_COMPILE_COMMANDS ON) - -# makes the current package a monorepo package -include(${CMAKE_CURRENT_SOURCE_DIR}/monorepoPackage.cmake) -set (VCPKG_OVERLAY_TRIPLETS ENV{VCPKG_OVERLAY_TRIPLETS}) - -set(CMAKE_TOOLCHAIN_FILE "$ENV{VCPKG_ROOT}/scripts/buildsystems/vcpkg.cmake") - -project(proto-rpc CXX) - -find_package(streamr-logger CONFIG REQUIRED) -find_package(streamr-eventemitter CONFIG REQUIRED) -find_package(streamr-utils CONFIG REQUIRED) -find_package(Protobuf REQUIRED) -find_package(Boost CONFIG REQUIRED) -find_package(magic_enum CONFIG REQUIRED) -find_package(folly REQUIRED) - -if(NOT TARGET Boost::uuid) - add_library(Boost::uuid INTERFACE IMPORTED) - set_target_properties(Boost::uuid PROPERTIES INTERFACE_INCLUDE_DIRECTORIES ${Boost_INCLUDE_DIRS}) -endif() - -add_library(streamr-proto-rpc) -set_property(TARGET streamr-proto-rpc PROPERTY CXX_STANDARD 26) -add_library(streamr::streamr-proto-rpc ALIAS streamr-proto-rpc) -target_include_directories(streamr-proto-rpc - PUBLIC $ - PUBLIC $ - PUBLIC $ - PUBLIC $ - ) - -target_sources(streamr-proto-rpc - PRIVATE ${CMAKE_CURRENT_SOURCE_DIR}/src/proto/packages/proto-rpc/protos/ProtoRpc.pb.cc) - -target_link_libraries(streamr-proto-rpc - INTERFACE streamr::streamr-logger - INTERFACE streamr::streamr-eventemitter - INTERFACE streamr::streamr-utils - PRIVATE protobuf::libprotobuf - PRIVATE protobuf::libprotoc - INTERFACE Boost::uuid - PRIVATE magic_enum::magic_enum - PUBLIC Folly::folly - ) - -export(TARGETS streamr-proto-rpc - NAMESPACE streamr:: - FILE streamr-proto-rpc-config-in.cmake) - -file(WRITE "${CMAKE_BINARY_DIR}/streamr-proto-rpc-config.cmake" - "list(APPEND CMAKE_PREFIX_PATH ${CMAKE_PREFIX_PATH})\n" - "find_package(streamr-logger CONFIG REQUIRED)\n" - "find_package(streamr-eventemitter CONFIG REQUIRED)\n" - "find_package(Protobuf REQUIRED)\n" - "find_package(Boost CONFIG REQUIRED)\n" - "find_package(magic_enum CONFIG REQUIRED)\n" - "find_package(folly REQUIRED)\n" - "if(NOT TARGET Boost::uuid)\n" - " add_library(Boost::uuid INTERFACE IMPORTED)\n" - " set_target_properties(Boost::uuid PROPERTIES INTERFACE_INCLUDE_DIRECTORIES ${Boost_INCLUDE_DIRS})\n" - "endif()\n" - "include(${CMAKE_BINARY_DIR}/streamr-proto-rpc-config-in.cmake)\n") - -if(NOT IOS) - add_executable(protobuf-streamr-plugin src/PluginCodeGeneratorMain.cpp include/streamr-proto-rpc/PluginCodeGenerator.hpp) - - target_include_directories( - protobuf-streamr-plugin PRIVATE $ - $) - - target_link_libraries(protobuf-streamr-plugin - PRIVATE protobuf::libprotobuf protobuf::libprotoc) - - - enable_testing() - find_package(GTest CONFIG REQUIRED) - - add_library(streamr-proto-rpc-test-main - src/CustomGtestMain.cpp - ) - - target_link_libraries(streamr-proto-rpc-test-main - PUBLIC Folly::folly - PUBLIC GTest::gtest - ) - - add_executable(streamr-proto-rpc-test-unit - test/proto/HelloRpc.pb.cc - test/proto/TestProtos.pb.cc - test/proto/WakeUpRpc.pb.cc - test/unit/ServerRegistryTest.cpp - test/unit/RpcCommunicatorTest.cpp - ) - - target_include_directories(streamr-proto-rpc-test-unit - PUBLIC $ - ) - - target_link_libraries(streamr-proto-rpc-test-unit - PUBLIC streamr-proto-rpc - PUBLIC GTest::gtest - PUBLIC streamr-proto-rpc-test-main - PUBLIC GTest::gmock - PUBLIC Folly::folly - ) - if (NOT (${VCPKG_TARGET_TRIPLET} MATCHES "android")) - include(GoogleTest) - gtest_discover_tests(streamr-proto-rpc-test-unit) - endif() - - add_executable(streamr-proto-rpc-test-integration - test/proto/HelloRpc.pb.cc - test/proto/HelloRpc.client.pb.h - test/proto/TestProtos.pb.cc - test/proto/WakeUpRpc.pb.cc - test/integration/ProtoRpcTest.cpp - ) - - target_include_directories(streamr-proto-rpc-test-integration - PUBLIC $ - ) - - target_link_libraries(streamr-proto-rpc-test-integration - PUBLIC streamr-proto-rpc - PUBLIC GTest::gtest - PUBLIC streamr-proto-rpc-test-main - PUBLIC GTest::gmock - PUBLIC Folly::folly - ) - if (NOT (${VCPKG_TARGET_TRIPLET} MATCHES "android")) - include(GoogleTest) - gtest_discover_tests(streamr-proto-rpc-test-integration) - endif() - - add_executable(streamr-proto-rpc-example-hello - examples/hello/hello.cpp - examples/hello/proto/HelloRpc.pb.cc - ) - - target_link_libraries(streamr-proto-rpc-example-hello - PUBLIC streamr-proto-rpc - PUBLIC Folly::folly - ) - - target_include_directories(streamr-proto-rpc-example-hello - PUBLIC $ - ) - - add_executable(streamr-proto-rpc-example-routed-hello - examples/routed-hello/routedhello.cpp - examples/routed-hello/proto/RoutedHelloRpc.pb.cc - ) - - target_link_libraries(streamr-proto-rpc-example-routed-hello - PUBLIC streamr-proto-rpc - PUBLIC Folly::folly - ) - - target_include_directories(streamr-proto-rpc-example-routed-hello - PUBLIC $ - ) -endif() diff --git a/packages/streamr-proto-rpc/include/streamr-proto-rpc/RpcCommunicator.hpp~ b/packages/streamr-proto-rpc/include/streamr-proto-rpc/RpcCommunicator.hpp~ deleted file mode 100644 index 3ce96813..00000000 --- a/packages/streamr-proto-rpc/include/streamr-proto-rpc/RpcCommunicator.hpp~ +++ /dev/null @@ -1,225 +0,0 @@ -#ifndef STREAMR_PROTO_RPC_RPC_COMMUNICATOR_HPP -#define STREAMR_PROTO_RPC_RPC_COMMUNICATOR_HPP - -#include -#include -#include -#include -#include -#include -#include -#include -#include -#include -#include -#include - -#include -#include "RpcCommunicatorClientApi.hpp" -#include "RpcCommunicatorServerApi.hpp" -#include "ServerRegistry.hpp" -#include "packages/proto-rpc/protos/ProtoRpc.pb.h" - -#include "streamr-logger/SLogger.hpp" - -namespace streamr::protorpc { - -using namespace std::chrono_literals; -using google::protobuf::Any; -using streamr::logger::SLogger; -using RpcMessage = ::protorpc::RpcMessage; -using RpcErrorType = ::protorpc::RpcErrorType; -using folly::coro::Task; - -constexpr std::chrono::milliseconds defaultRpcRequestTimeout = 5000ms; - -// NOLINTBEGIN -enum class StatusCode { OK, STOPPED, DEADLINE_EXCEEDED, SERVER_ERROR }; -// NOLINTEND - -struct RpcCommunicatorOptions { - std::chrono::milliseconds rpcRequestTimeout; -}; - -template -class RpcCommunicator { -public: - using OutgoingMessageCallbackType = std::function; - -private: - RpcCommunicatorClientApi - mRpcCommunicatorClientApi; - RpcCommunicatorServerApi - mRpcCommunicatorServerApi; - -public: - using RequestId = - RpcCommunicatorClientApi:: - RequestId; - using OngoingRequestBase = - RpcCommunicatorClientApi:: - OngoingRequestBase; - - explicit RpcCommunicator( - std::optional options = std::nullopt) - : mRpcCommunicatorClientApi( - options.has_value() ? options.value().rpcRequestTimeout - : defaultRpcRequestTimeout) {} - - // Messaging API - - /** - * @brief Handle a message incoming from the network - * - * @param message message received from the network - * @param callContext call context such as routing information - */ - - void handleIncomingMessage( - const RpcMessage& message, const CallContextType& callContext) { - this->onIncomingMessage(message, callContext); - } - - /** - * @brief Set a callback for sending outgoing messages to the network - * - * @param callback callback to be called when a message is ready to be sent. - * The callback may throw exeptions which will be forwarded through the - * RpcCommunicator - * - */ - - void setOutgoingMessageCallback( - const OutgoingMessageCallbackType& callback) { - mRpcCommunicatorClientApi.setOutgoingMessageCallback(callback); - mRpcCommunicatorServerApi.setOutgoingMessageCallback(callback); - } - - // Client-side API - - /** - * @brief [A method to be called by auto-generated clients] Make a RPC - * request to a remote method - * - * @param methodName name of the method to be called - * @param methodParam parameter to be passed to the method - * @param callContext call context such as routing information - * @return Task - */ - - template - Task request( - const std::string& methodName, - const RequestType& methodParam, - const CallContextType& callContext, - std::optional timeout = std::nullopt) { - return mRpcCommunicatorClientApi - .template request( - methodName, methodParam, callContext, timeout); - } - - /** - * @brief [A method to be called by auto-generated clients] Make a - * remote RPC notification - * - * @param notificationName name of the notification to be called - * @param notificationParam parameter to be passed to the notification - * @param callContext call context such as routing information - * @return Task - */ - - template - Task notify( - const std::string_view notificationName, - const RequestType& notificationParam, - const CallContextType& callContext, - std::optional timeout = std::nullopt) { - return mRpcCommunicatorClientApi.template notify( - notificationName, notificationParam, callContext, timeout); - } - - // Server-side API - - /** - * @brief Register a method to be called by remote clients - * - * @param name name of the method - * @param fn function to be registered - * @param options options for the method - */ - - template - requires std::is_assignable_v< - std::function, - F> - void registerRpcMethod( - const std::string& name, const F& fn, MethodOptions options = {}) { - mRpcCommunicatorServerApi - .template registerRpcMethod( - name, fn, options); - } - - /** - * @brief Register a notification to be called by remote clients - * - * @param name name of the notification - * @param fn function to be registered - * @param options options for the notification - */ - - template - requires std::is_assignable_v< - std::function, - F> - void registerRpcNotification( - const std::string& name, const F& fn, MethodOptions options = {}) { - mRpcCommunicatorServerApi - .template registerRpcNotification( - name, fn, options); - } - -private: - void onIncomingMessage( - const RpcMessage& rpcMessage, const CallContextType& callContext) { - SLogger::trace("onIncomingMessage", rpcMessage.DebugString()); - - const auto& header = rpcMessage.header(); - - SLogger::trace("onIncomingMessage() requestId", rpcMessage.requestid()); - - if (header.find("response") != header.end()) { - mRpcCommunicatorClientApi.onIncomingMessage( - rpcMessage, callContext); - } else if ( - header.find("request") != header.end() && - header.find("method") != header.end()) { - mRpcCommunicatorServerApi.onIncomingMessage( - rpcMessage, callContext); - } else { - SLogger::debug( - "onIncomingMessage() message is not a valid request or response"); - } - } - -protected: - using OngoingRequestPredicate = - RpcCommunicatorClientApi:: - OngoingRequestPredicate; - - [[nodiscard]] std::vector - getOngoingRequestIdsFulfillingPredicate( - const OngoingRequestPredicate& predicate) { - return mRpcCommunicatorClientApi - .getOngoingRequestIdsFulfillingPredicate(predicate); - } - - void handleClientError( - const RequestId& requestId, const RpcException& error) { - mRpcCommunicatorClientApi.handleClientError(requestId, error); - } -}; - -} // namespace streamr::protorpc - -#endif // STREAMR_PROTO_RPC_RPC_COMMUNICATOR_HPP \ No newline at end of file diff --git a/packages/streamr-proto-rpc/test/unit/RpcCommunicatorTest.cpp~ b/packages/streamr-proto-rpc/test/unit/RpcCommunicatorTest.cpp~ deleted file mode 100644 index a84d5d17..00000000 --- a/packages/streamr-proto-rpc/test/unit/RpcCommunicatorTest.cpp~ +++ /dev/null @@ -1,472 +0,0 @@ -#include -#include -#include -#include -#include -#include -#include -#include -#include -#include -#include -#include -#include -#include -#include -#include -#include -#include -#include "HelloRpc.pb.h" -#include "streamr-proto-rpc/Errors.hpp" -#include "streamr-proto-rpc/ProtoCallContext.hpp" -#include "streamr-proto-rpc/ServerRegistry.hpp" - -namespace streamr::protorpc { - -using namespace std::chrono_literals; -using RpcCommunicatorType = RpcCommunicator; -class RpcCommunicatorTest : public ::testing::Test { -public: - RpcCommunicatorTest() - : executor(folly::CPUThreadPoolExecutor(threadPoolSize)) {} // NOLINT - - ~RpcCommunicatorTest() override { - SLogger::warn("Deleting executor of RpcCommunicatorTest"); - SLogger::warn("RpcCommunicatorTypeTest executor deleted"); - } - -protected: - RpcCommunicatorType communicator1; // NOLINT - RpcCommunicatorType communicator2; // NOLINT - folly::CPUThreadPoolExecutor executor; // NOLINT - - void SetUp() override {} - - void setCallbacks(bool isMethod = true) { - communicator2.setOutgoingMessageCallback( - [this]( - const RpcMessage& message, - const std::string& /* requestId */, - const ProtoCallContext& /* context */) -> void { - SLogger::info("onOutgoingMessageCallback()"); - communicator1.handleIncomingMessage( - message, ProtoCallContext()); - }); - if (isMethod) { - communicator1.setOutgoingMessageCallback( - [this]( - const RpcMessage& message, - const std::string& /* requestId */, - const ProtoCallContext& /* context */) -> void { - SLogger::info("onOutgoingMessageCallback()"); - communicator2.handleIncomingMessage( - message, ProtoCallContext()); - }); - } - } -}; - -void registerTestRcpMethod(RpcCommunicatorType& communicator) { - communicator.registerRpcMethod( - "testFunction", - [](const HelloRequest& request, - const ProtoCallContext& /* context */) -> HelloResponse { - HelloResponse response; - SLogger::info("testFunction() request.myname():", request.myname()); - response.set_greeting("Hello, " + request.myname()); - return response; - }); -} - -template -void registerThrowingRcpMethod(RpcCommunicatorType& communicator) { - communicator.registerRpcMethod( - "testFunction", - [](const HelloRequest& /* request */, - const ProtoCallContext& /* context */) -> HelloResponse { - throw T("TestException"); - }); -} - -void registerThrowingTestRcpMethodUnknown(RpcCommunicatorType& communicator) { - const int unknownThrow = 42; - communicator.registerRpcMethod( - "testFunction", - [](const HelloRequest& /* request */, - const ProtoCallContext& /* context */) -> HelloResponse { - throw unknownThrow; // NOLINT - }); -} -/* -void setOutgoingCallback(RpcCommunicatorType& sender, RpcCommunicatorType& -receiver) { sender.setOutgoingMessageCallback( - [&receiver]( - const RpcMessage& message, - const std::string& requestId , - const ProtoCallContext& context ) -> void { - SLogger::info("onOutgoingMessageCallback()"); - receiver.handleIncomingMessage(message, ProtoCallContext()); - }); -} - -void setOutgoingCallbacks( - RpcCommunicatorType& communicator1, RpcCommunicatorType& communicator2) { - setOutgoingCallback(communicator2, communicator1); - setOutgoingCallback(communicator1, communicator2); -} -*/ -template -void setOutgoingCallbackWithException(RpcCommunicatorType& sender) { - sender.setOutgoingMessageCallback( - [](const RpcMessage& /* message */, - const std::string& /* requestId */, - const ProtoCallContext& /* context */) -> void { - SLogger::info("onOutgoingMessageCallback() throws"); - throw T("TestException"); - }); -} - -void setOutgoingCallbackWithThrownUnknown(RpcCommunicatorType& sender) { - const int unknownThrow = 42; - sender.setOutgoingMessageCallback( - [](const RpcMessage& /* message */, - const std::string& /* requestId */, - const ProtoCallContext& /* context */) -> void { - SLogger::info("onOutgoingMessageCallback() throws"); - throw unknownThrow; // NOLINT - }); -} - -auto sendHelloRequest( - RpcCommunicatorType& sender, folly::CPUThreadPoolExecutor* executor) { - HelloRequest request; - request.set_myname("Test"); - return folly::coro::blockingWait( - sender - .request( - "testFunction", request, ProtoCallContext()) - .scheduleOn(executor)); -} - -auto sendHelloNotification( - RpcCommunicatorType& sender, folly::CPUThreadPoolExecutor* executor) { - HelloRequest request; - request.set_myname("Test"); - folly::coro::blockingWait( - sender.notify("testFunction", request, ProtoCallContext()) - .scheduleOn(executor)); -} - -void verifyClientError( - const RpcClientError& ex, - const ErrorCode expectedErrorCode, - const std::string& expectedOriginalErrorInfo) { - SLogger::info("Caught RpcClientError", ex.what()); - EXPECT_EQ(ex.code, expectedErrorCode); - EXPECT_TRUE(ex.originalErrorInfo.has_value()); - EXPECT_EQ(ex.originalErrorInfo.value(), expectedOriginalErrorInfo); -} - -void registerSleepingTestRcpMethod(RpcCommunicatorType& communicator) { - communicator.registerRpcMethod( - "testFunction", - [](const HelloRequest& request, - const ProtoCallContext& /* context */) -> HelloResponse { - SLogger::info("TestSleepingRpcMethod sleeping 1s"); - std::this_thread::sleep_for(std::chrono::seconds(1)); // NOLINT - HelloResponse response; - response.set_greeting("Hello, " + request.myname()); - return response; - }); -} - -TEST_F(RpcCommunicatorTest, TestCanRegisterRpcMethod) { - RpcCommunicatorType communicator; - registerTestRcpMethod(communicator); - EXPECT_EQ(true, true); -} - -TEST_F(RpcCommunicatorTest, TestCanMakeRpcCall) { - registerTestRcpMethod(communicator1); - setCallbacks(); - auto result = sendHelloRequest(communicator2, &executor); - EXPECT_EQ("Hello, Test", result.greeting()); -} - -TEST_F(RpcCommunicatorTest, TestrequestClientThrowsRuntimeError) { - registerTestRcpMethod(communicator1); - setOutgoingCallbackWithException(communicator2); - try { - sendHelloRequest(communicator2, &executor); - EXPECT_TRUE(false); - } catch (const RpcClientError& ex) { - verifyClientError(ex, ErrorCode::RPC_CLIENT_ERROR, "TestException"); - } catch (...) { - EXPECT_TRUE(false); - } -} - -TEST_F(RpcCommunicatorTest, TestrequestClientThrowsFailedToParse) { - registerTestRcpMethod(communicator1); - setOutgoingCallbackWithException(communicator2); - try { - sendHelloRequest(communicator2, &executor); - EXPECT_TRUE(false); - } catch (const RpcClientError& ex) { - verifyClientError( - ex, - ErrorCode::RPC_CLIENT_ERROR, - "code: FAILED_TO_PARSE, message: TestException"); - } catch (...) { - EXPECT_TRUE(false); - } -} - -TEST_F(RpcCommunicatorTest, TestrequestClientThrowsUnknownError) { - registerTestRcpMethod(communicator1); - setOutgoingCallbackWithThrownUnknown(communicator2); - try { - sendHelloRequest(communicator2, &executor); - EXPECT_TRUE(false); - } catch (...) { - EXPECT_TRUE(true); - } -} - -TEST_F(RpcCommunicatorTest, TestrequestServerThrowsRuntimeError) { - registerThrowingRcpMethod(communicator1); - setCallbacks(); - try { - sendHelloRequest(communicator2, &executor); - EXPECT_TRUE(false); - } catch (const RpcServerError& ex) { - EXPECT_EQ(ex.code, ErrorCode::RPC_SERVER_ERROR); - EXPECT_TRUE( - ex.errorClassName.find("runtime_error") != std::string::npos); - EXPECT_EQ(ex.message, "TestException"); - } catch (...) { - EXPECT_TRUE(false); - } -} - -TEST_F(RpcCommunicatorTest, TestrequestServerThrowsUnknownRpcMethod) { - registerThrowingRcpMethod(communicator1); - setCallbacks(); - try { - sendHelloRequest(communicator2, &executor); - EXPECT_TRUE(false); - } catch (const UnknownRpcMethod& ex) { - EXPECT_EQ(ex.code, ErrorCode::UNKNOWN_RPC_METHOD); - } catch (...) { - EXPECT_TRUE(false); - } -} - -TEST_F(RpcCommunicatorTest, TestrequestServerThrowsRcpTimeout) { - registerThrowingRcpMethod(communicator1); - setCallbacks(); - try { - sendHelloRequest(communicator2, &executor); - EXPECT_TRUE(false); - } catch (const RpcTimeout& ex) { - EXPECT_EQ(ex.code, ErrorCode::RPC_TIMEOUT); - } catch (...) { - EXPECT_TRUE(false); - } -} - -TEST_F(RpcCommunicatorTest, TestrequestServerThrowsFailedToParse) { - registerThrowingRcpMethod(communicator1); - setCallbacks(); - try { - sendHelloRequest(communicator2, &executor); - EXPECT_TRUE(false); - } catch (const RpcServerError& ex) { - EXPECT_EQ(ex.code, ErrorCode::RPC_SERVER_ERROR); - EXPECT_EQ(ex.errorCode, "FAILED_TO_PARSE"); - EXPECT_TRUE( - ex.errorClassName.find("FailedToParse") != std::string::npos); - } catch (...) { - EXPECT_TRUE(false); - } -} - -TEST_F(RpcCommunicatorTest, TestrequestServerThrowsUnknown) { - registerThrowingTestRcpMethodUnknown(communicator1); - setCallbacks(); - try { - sendHelloRequest(communicator2, &executor); - EXPECT_TRUE(false); - } catch (...) { - EXPECT_TRUE(true); - } -} - -TEST_F(RpcCommunicatorTest, TestCannotify) { - std::string requestMsg; - communicator1.registerRpcNotification( - "testFunction", - [&requestMsg]( - const HelloRequest& request, const ProtoCallContext& /* context */) - -> void { requestMsg = request.DebugString(); }); - setCallbacks(false); - sendHelloNotification(communicator2, &executor); - EXPECT_EQ(requestMsg, "myName: \"Test\"\n"); -} - -TEST_F(RpcCommunicatorTest, TestnotifyClientThrowsRuntimeError) { - setOutgoingCallbackWithException(communicator2); - try { - sendHelloNotification(communicator2, &executor); - EXPECT_TRUE(false); - } catch (const RpcClientError& ex) { - verifyClientError(ex, ErrorCode::RPC_CLIENT_ERROR, "TestException"); - } catch (...) { - EXPECT_TRUE(false); - } -} - -TEST_F(RpcCommunicatorTest, TestnotifyClientThrowsFailedToParse) { - setOutgoingCallbackWithException(communicator2); - try { - sendHelloNotification(communicator2, &executor); - EXPECT_TRUE(false); - } catch (const RpcClientError& ex) { - verifyClientError( - ex, - ErrorCode::RPC_CLIENT_ERROR, - "code: FAILED_TO_PARSE, message: TestException"); - } catch (...) { - EXPECT_TRUE(false); - } -} - -TEST_F(RpcCommunicatorTest, TestRpcTimeoutOnClientSide) { - registerTestRcpMethod(communicator1); - - communicator2.setOutgoingMessageCallback( - [&communicator1 = this->communicator1]( - const RpcMessage& /* message */, - const std::string& /* requestId */, - const ProtoCallContext& /* context */) -> void { - SLogger::info("setOutgoingMessageCallback() sleeping 1s"); - std::this_thread::sleep_for(std::chrono::seconds(1)); // NOLINT - }); - - HelloRequest request; - request.set_myname("Test"); - - try { - folly::coro::blockingWait(communicator2 - .request( - "testFunction", - request, - ProtoCallContext(), - 50ms) // NOLINT - .scheduleOn(&executor)); - // Test fails here - EXPECT_TRUE(false); - } catch (const RpcTimeout& ex) { - SLogger::info("TestRpcTimeout caught RpcTimeout", ex.what()); - } catch (const std::exception& ex) { - SLogger::info( - "TestRpcTimeoutOnClientSide caught unknown exception", ex.what()); - EXPECT_TRUE(false); - } -} - -TEST_F(RpcCommunicatorTest, TestRpcTimeoutOnServerSide) { - registerSleepingTestRcpMethod(communicator1); - - std::shared_ptr thread; - - communicator2.setOutgoingMessageCallback( - [&communicator1 = this->communicator1, &thread]( - const RpcMessage& message, - const std::string& /* requestId */, - const ProtoCallContext& context) -> void { - thread = std::make_shared( - [&communicator1, message, context]() { - SLogger::info("Starting thread for server"); - communicator1.handleIncomingMessage(message, context); - }); - }); - communicator1.setOutgoingMessageCallback( - [&communicator2 = this->communicator2]( - const RpcMessage& message, - const std::string& /* requestId */, - const ProtoCallContext& context) -> void { - communicator2.handleIncomingMessage(message, context); - }); - - HelloRequest request; - request.set_myname("Test"); - try { - auto result = - folly::coro::blockingWait(communicator2 - .request( - "testFunction", - request, - ProtoCallContext(), - 50ms) // NOLINT - .scheduleOn(&executor)); - EXPECT_EQ(true, false); - } catch (const RpcTimeout& ex) { - SLogger::info( - "TestRpcTimeoutOnServerSide caught RpcTimeout", ex.what()); - } catch (...) { - SLogger::info("TestRpcTimeoutOnServerSide caught unknown exception"); - EXPECT_EQ(true, false); - } - thread->join(); -} - -TEST_F(RpcCommunicatorTest, TestRpcTimeoutOnClientSideForNotification) { - std::string requestMsg; - - communicator1.registerRpcNotification( - "testFunction", - [&requestMsg]( - const HelloRequest& request, const ProtoCallContext& /* context */) - -> void { requestMsg = request.DebugString(); }); - - HelloRequest request; - request.set_myname("Test"); - - communicator2.setOutgoingMessageCallback( - [&communicator1 = this->communicator1]( - const RpcMessage& /* message */, - const std::string& /* requestId */, - const ProtoCallContext& /* context */) -> void { - std::cout << "setOutgoingMessageCallback() sleeping 1s, thread id: " - << std::this_thread::get_id() << "\n"; - std::this_thread::sleep_for(std::chrono::seconds(1)); // NOLINT - }); - try { - // auto str = std::this_thread::get_id(); - - std::cout << "Calling notify() from thread id: " - << std::this_thread::get_id() << "\n"; - folly::coro::blockingWait(communicator2 - .notify( - "testFunction", - request, - ProtoCallContext(), - 50ms) // NOLINT - .scheduleOn(&executor)); - // Test fails here - EXPECT_TRUE(false); - } catch (const RpcTimeout& ex) { - std::cout - << "TestRpcTimeoutOnClientSideForNotification caught RpcTimeout in thread id: " - << std::this_thread::get_id() << "\n"; - } catch (const std::exception& ex) { - std::cout - << "TestRpcTimeoutOnClientSideForNotification caught unknown exception in thread id: " - << std::this_thread::get_id() << "\n"; - EXPECT_TRUE(false); - EXPECT_TRUE(false); - } -} - -} // namespace streamr::protorpc