Skip to content
Merged
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
5 changes: 5 additions & 0 deletions packages/streamr-dht/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -172,6 +172,8 @@ if(NOT IOS AND STREAMR_MODULES_SUPPORTED)
test/unit/WebsocketClientConnectorRpcRemoteTest.cpp
test/unit/IdentifiersTest.cpp
test/unit/CreatePeerDescriptorTest.cpp
test/unit/WebrtcConnectionTest.cpp
test/unit/WebrtcConnectorTest.cpp
test/unit/ConnectivityRequestHandlerTest.cpp
test/unit/ConnectionLockRpcRemoteTest.cpp
test/unit/ConnectorFacadeTest.cpp
Expand Down Expand Up @@ -264,6 +266,9 @@ if(NOT IOS AND STREAMR_MODULES_SUPPORTED)
test/integration/KademliaCorrectnessTest.cpp
test/integration/ConnectivityCheckingTest.cpp
test/integration/WebsocketConnectionManagementTest.cpp
test/integration/WebrtcConnectorRpcTest.cpp
test/integration/WebrtcConnectionManagementTest.cpp
test/integration/RpcConnectionsOverWebrtcTest.cpp
)

target_include_directories(streamr-dht-test-integration
Expand Down
9 changes: 8 additions & 1 deletion packages/streamr-dht/lint.sh
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,13 @@ echo "Running clangd-tidy on $FILES"
# CodeLocation/MakeAndRegisterTestInfo reported ambiguous — like the other
# excluded files; the compiler builds and runs it on every platform and
# clang-format still checks it.)
# WebrtcConnectionTest.cpp, WebrtcConnectorTest.cpp and
# RpcConnectionsOverWebrtcTest.cpp (phase B2) trip the same std-type
# unification false positive (gtest CodeLocation/MakeAndRegisterTestInfo
# ambiguous, protobuf member calls ambiguous) through the WebrtcConnection /
# ConnectionManager import sets; the compiler builds and runs all three
# (unit 181/181, integration green). The other two B2 tests
# (all five B2 test files trip it via the protobuf/gtest ambiguity.)
# CreatePeerDescriptorTest.cpp, ConnectivityRequestHandlerTest.cpp and
# WebsocketConnectionManagementTest.cpp (phase B1) trip the same std-type
# unification false positive (std::string method calls on protobuf getters
Expand All @@ -84,7 +91,7 @@ echo "Running clangd-tidy on $FILES"
# clang-format still checks them. (ConnectivityCheckingTest.cpp imports
# the same cluster but does NOT trip it.)
HEAVY_TIDY_FILES='./test/integration/MockLayer1Layer0Test.cpp ./test/integration/Layer1ScaleTest.cpp ./test/integration/DhtNodeExternalApiTest.cpp ./test/integration/KademliaCorrectnessTest.cpp'
TIDY_FILES=$(echo "$FILES" | tr ' ' '\n' | grep -v 'test/integration/ConnectionLockingTest.cpp' | grep -v 'test/unit/PendingConnectionTest.cpp' | grep -v 'test/unit/SimulatorTest.cpp' | grep -v 'test/integration/SimultaneousConnectionsTest.cpp' | grep -v 'test/integration/ConnectionManagerIntegrationTest.cpp' | grep -v 'test/unit/DhtNodeRpcLocalTest.cpp' | grep -v 'test/integration/DhtNodeRpcRemoteTest.cpp' | grep -v 'test/unit/RoutingSessionTest.cpp' | grep -v 'test/unit/RecursiveOperationSessionTest.cpp' | grep -v 'test/unit/StoreRpcLocalTest.cpp' | grep -v 'test/unit/StoreManagerTest.cpp' | grep -v 'test/integration/MultipleEntryPointJoiningTest.cpp' | grep -v 'test/utils/DhtNodeTestUtils.hpp' | grep -v 'test/integration/DhtNodeTest.cpp' | grep -v 'test/integration/MockLayer1Layer0Test.cpp' | grep -v 'test/integration/Layer1ScaleTest.cpp' | grep -v 'test/integration/DhtNodeExternalApiTest.cpp' | grep -v 'test/integration/KademliaCorrectnessTest.cpp' | grep -v 'test/unit/CreatePeerDescriptorTest.cpp' | grep -v 'test/unit/ConnectivityRequestHandlerTest.cpp' | grep -v 'test/integration/WebsocketConnectionManagementTest.cpp' | tr '\n' ' ')
TIDY_FILES=$(echo "$FILES" | tr ' ' '\n' | grep -v 'test/integration/ConnectionLockingTest.cpp' | grep -v 'test/unit/PendingConnectionTest.cpp' | grep -v 'test/unit/SimulatorTest.cpp' | grep -v 'test/integration/SimultaneousConnectionsTest.cpp' | grep -v 'test/integration/ConnectionManagerIntegrationTest.cpp' | grep -v 'test/unit/DhtNodeRpcLocalTest.cpp' | grep -v 'test/integration/DhtNodeRpcRemoteTest.cpp' | grep -v 'test/unit/RoutingSessionTest.cpp' | grep -v 'test/unit/RecursiveOperationSessionTest.cpp' | grep -v 'test/unit/StoreRpcLocalTest.cpp' | grep -v 'test/unit/StoreManagerTest.cpp' | grep -v 'test/integration/MultipleEntryPointJoiningTest.cpp' | grep -v 'test/utils/DhtNodeTestUtils.hpp' | grep -v 'test/integration/DhtNodeTest.cpp' | grep -v 'test/integration/MockLayer1Layer0Test.cpp' | grep -v 'test/integration/Layer1ScaleTest.cpp' | grep -v 'test/integration/DhtNodeExternalApiTest.cpp' | grep -v 'test/integration/KademliaCorrectnessTest.cpp' | grep -v 'test/unit/CreatePeerDescriptorTest.cpp' | grep -v 'test/unit/ConnectivityRequestHandlerTest.cpp' | grep -v 'test/integration/WebsocketConnectionManagementTest.cpp' | grep -v 'test/unit/WebrtcConnectionTest.cpp' | grep -v 'test/unit/WebrtcConnectorTest.cpp' | grep -v 'test/integration/RpcConnectionsOverWebrtcTest.cpp' | grep -v 'test/integration/WebrtcConnectorRpcTest.cpp' | grep -v 'test/integration/WebrtcConnectionManagementTest.cpp' | tr '\n' ' ')
# Chunked: one clangd process accumulates source-location space across the
# files it serves and never releases it; with the phase-A8 module graph a
# single process no longer survives the whole list (mid-batch the affected
Expand Down
7 changes: 7 additions & 0 deletions packages/streamr-dht/modules/connection/Connection.cppm
Original file line number Diff line number Diff line change
Expand Up @@ -76,6 +76,13 @@ public:
~Connection() override { SLogger::trace("~Connection()"); }

[[nodiscard]] ConnectionID getConnectionID() const { return mID; }

// TS IConnection.connectionId is publicly assignable; used by the WebRTC
// answerer to adopt the offerer's connection id (WebrtcConnectorRpcLocal
// rtcOffer).
void setConnectionId(ConnectionID connectionId) {
this->mID = std::move(connectionId);
}
[[nodiscard]] ConnectionType getConnectionType() const { return mType; }
[[nodiscard]] std::string getConnectionTypeString() const {
return std::string(magic_enum::enum_name(mType));
Expand Down
45 changes: 34 additions & 11 deletions packages/streamr-dht/modules/connection/ConnectorFacade.cppm
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,8 @@ import streamr.dht.ListeningRpcCommunicator;
import streamr.dht.PendingConnection;
import streamr.dht.PortRange;
import streamr.dht.Transport;
import streamr.dht.WebrtcConnector;
import streamr.dht.webrtcTypes;
import streamr.dht.WebsocketClientConnector;
import streamr.dht.WebsocketServerConnector;

Expand All @@ -40,6 +42,9 @@ using namespace std::chrono_literals;
using ::dht::ConnectivityResponse;
using streamr::dht::connection::IPendingConnection;
using streamr::dht::connection::PendingConnection;
using streamr::dht::connection::webrtc::IceServer;
using streamr::dht::connection::webrtc::WebrtcConnector;
using streamr::dht::connection::webrtc::WebrtcConnectorOptions;
using streamr::dht::connection::websocket::WebsocketClientConnector;
using streamr::dht::connection::websocket::WebsocketClientConnectorOptions;
using streamr::dht::connection::websocket::WebsocketServerConnector;
Expand Down Expand Up @@ -71,12 +76,12 @@ struct DefaultConnectorFacadeOptions {
std::optional<std::string> websocketHost = std::nullopt;
std::optional<PortRange> websocketPortRange = std::nullopt;
std::optional<std::vector<PeerDescriptor>> entryPoints = std::nullopt;
// std::vector<IceServer> iceServers;
// bool webrtcAllowPrivateAddresses;
// int webrtcDatachannelBufferThresholdLow;
// int webrtcDatachannelBufferThresholdHigh;
// std::optional<std::string> externalIp;
// PortRange webrtcPortRange;
std::vector<IceServer> iceServers = {};
std::optional<bool> webrtcAllowPrivateAddresses = std::nullopt;
std::optional<size_t> webrtcDatachannelBufferThresholdLow = std::nullopt;
std::optional<size_t> webrtcDatachannelBufferThresholdHigh = std::nullopt;
std::optional<std::string> externalIp = std::nullopt;
std::optional<PortRange> webrtcPortRange = std::nullopt;
std::optional<size_t> maxMessageSize;
// TlsCertificate tlsCertificate;
// bool websocketServerEnableTls;
Expand All @@ -97,11 +102,13 @@ private:
ListeningRpcCommunicator websocketConnectorRpcCommunicator;
std::unique_ptr<WebsocketClientConnector> websocketClientConnector;
std::unique_ptr<WebsocketServerConnector> websocketServerConnector;
std::unique_ptr<WebrtcConnector> webrtcConnector;

void setLocalPeerDescriptor(const PeerDescriptor& peerDescriptor) {
this->localPeerDescriptor = peerDescriptor;
this->websocketClientConnector->setLocalPeerDescriptor(peerDescriptor);
this->websocketServerConnector->setLocalPeerDescriptor(peerDescriptor);
this->webrtcConnector->setLocalPeerDescriptor(peerDescriptor);
}

public:
Expand Down Expand Up @@ -151,6 +158,21 @@ public:
std::make_unique<WebsocketServerConnector>(
std::move(webSocketServerConnectorOptions));

this->webrtcConnector =
std::make_unique<WebrtcConnector>(WebrtcConnectorOptions{
.onNewConnection = onNewConnection,
.transport = this->options.transport,
.iceServers = this->options.iceServers,
.allowPrivateAddresses =
this->options.webrtcAllowPrivateAddresses,
.bufferThresholdLow =
this->options.webrtcDatachannelBufferThresholdLow,
.bufferThresholdHigh =
this->options.webrtcDatachannelBufferThresholdHigh,
.maxMessageSize = this->options.maxMessageSize,
.externalIp = this->options.externalIp,
.portRange = this->options.webrtcPortRange});

this->websocketServerConnector->start();
const auto connectivityResponse =
this->websocketServerConnector->checkConnectivity(false);
Expand Down Expand Up @@ -178,8 +200,7 @@ public:
peerDescriptor)) {
return this->websocketServerConnector->connect(peerDescriptor);
}
// TS falls through to the WebrtcConnector here (milestone B2).
return nullptr;
return this->webrtcConnector->connect(peerDescriptor, false);
}

std::shared_ptr<IPendingConnection> createConnection(
Expand All @@ -194,8 +215,9 @@ public:
peerDescriptor)) {
return this->websocketServerConnector->connect(peerDescriptor);
}
// TS falls through to the WebrtcConnector here (milestone B2).
return nullptr;
// The TS PendingConnection carries no per-call error callback on the
// WebRTC path either.
return this->webrtcConnector->connect(peerDescriptor, false);
}

PeerDescriptor getLocalPeerDescriptor() const override {
Expand All @@ -205,8 +227,9 @@ public:
void stop() override {
SLogger::info("DefaultConnectorFacade::stop start");
this->websocketConnectorRpcCommunicator.destroy();
this->websocketClientConnector->destroy();
this->websocketServerConnector->destroy();
this->websocketClientConnector->destroy();
this->webrtcConnector->stop();
SLogger::info("DefaultConnectorFacade::stop end");
}
};
Expand Down
Loading
Loading