From 4bdf9658cc21940c37dd7ba4e9dc5feafe5ef50f Mon Sep 17 00:00:00 2001 From: Petri Savolainen Date: Thu, 2 Jul 2026 07:28:28 +0300 Subject: [PATCH 1/2] Phase 1.0: repo baseline fixes - Remove 5 accidentally committed editor backup files (*~) and ignore the pattern in .gitignore. - Pin wrappers/go submodule to reachable commit 55f17ba (v2.0.0). The previously recorded commit cdb68c95 no longer exists on the remote, which made `git submodule update --init --recursive` fail for every fresh clone (it also aborts checkout of the remaining submodules, leaving vcpkg at whatever HEAD the clone fetched instead of the pinned baseline commit). Part of the toolchain + C++ modules modernization, Phase 1.0. Co-Authored-By: Claude Fable 5 --- .gitignore | 1 + .../streamr-json/test/unit/toJsonTest.cpp~ | 367 -------------- .../CMakeLists.txt~ | 93 ---- .../streamr-libstreamrproxyclient/wrappers/go | 2 +- packages/streamr-proto-rpc/CMakeLists.txt~ | 170 ------- .../streamr-proto-rpc/RpcCommunicator.hpp~ | 225 --------- .../test/unit/RpcCommunicatorTest.cpp~ | 472 ------------------ 7 files changed, 2 insertions(+), 1328 deletions(-) delete mode 100644 packages/streamr-json/test/unit/toJsonTest.cpp~ delete mode 100644 packages/streamr-libstreamrproxyclient/CMakeLists.txt~ delete mode 100644 packages/streamr-proto-rpc/CMakeLists.txt~ delete mode 100644 packages/streamr-proto-rpc/include/streamr-proto-rpc/RpcCommunicator.hpp~ delete mode 100644 packages/streamr-proto-rpc/test/unit/RpcCommunicatorTest.cpp~ 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 From 0c376baa9bfbdefdc3a89b229bb2e11f2bc9bf5f Mon Sep 17 00:00:00 2001 From: Petri Savolainen Date: Thu, 2 Jul 2026 07:37:36 +0300 Subject: [PATCH 2/2] CI: disable matrix fail-fast in validate workflow One platform's failure was cancelling all other matrix legs, which hides the per-platform status this modernization needs to track (macOS legs are expected red until the Phase 1.2 compiler upgrade; the Linux legs are the old-toolchain baseline gate). Co-Authored-By: Claude Fable 5 --- .github/workflows/validate.yml | 4 ++++ 1 file changed, 4 insertions(+) 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 }}