diff --git a/MODERNIZATION.md b/MODERNIZATION.md index 2d2711b3..6cc9e18b 100644 --- a/MODERNIZATION.md +++ b/MODERNIZATION.md @@ -170,6 +170,15 @@ Android NDK r28+. `--repeat until-pass:2 --timeout 300`) keep hangs bounded and retries honest; the tests' timing/port assumptions need their own fix. This debt predates the modernization and was exposed by introducing macOS CI at all. + **RESOLVED post-modernization:** the flakes were two library races in + `WebsocketConnection`, not test timing/port assumptions — (1) incoming + frames emitted as Data events before the application could register a + listener (message silently dropped; fixed by deferring the rtc message + callback to `startReceiving()`), and (2) a lock-order inversion between + `mMutex` and rtc's callback mutex in `close()`/`destroy()` (teardown + deadlock; fixed by calling into rtc outside `mMutex`). Reproduced 3/200 + resp. 2/100 locally before the fix; 0 failures in 925 stress runs + (Debug + Release) after. - **Gate**: build/test green macOS + Linux, **and iOS cross-build + `iostest.sh` green — the compiler's output must stay compatible with the device's fixed libc++ runtime**. diff --git a/packages/streamr-dht/include/streamr-dht/connection/websocket/WebsocketClientConnection.hpp b/packages/streamr-dht/include/streamr-dht/connection/websocket/WebsocketClientConnection.hpp index 1fa69260..cc55dd25 100644 --- a/packages/streamr-dht/include/streamr-dht/connection/websocket/WebsocketClientConnection.hpp +++ b/packages/streamr-dht/include/streamr-dht/connection/websocket/WebsocketClientConnection.hpp @@ -41,6 +41,10 @@ class WebsocketClientConnection : public WebsocketConnection { auto socket = std::make_shared(webSocketConfig); setSocket(socket); + // The client's listeners are registered before connect(), so the + // message callback can be attached right away — before open(), as + // it always was (setSocket() no longer attaches it itself). + startReceiving(); SLogger::trace("socket created"); socket->open(address); SLogger::trace("connect() mSocket->open() called"); diff --git a/packages/streamr-dht/include/streamr-dht/connection/websocket/WebsocketConnection.hpp b/packages/streamr-dht/include/streamr-dht/connection/websocket/WebsocketConnection.hpp index 5389643e..9000d8ef 100644 --- a/packages/streamr-dht/include/streamr-dht/connection/websocket/WebsocketConnection.hpp +++ b/packages/streamr-dht/include/streamr-dht/connection/websocket/WebsocketConnection.hpp @@ -29,6 +29,7 @@ class WebsocketConnection : public Connection, public EnableSharedFromThis { protected: std::shared_ptr mSocket; // NOLINT std::atomic mDestroyed{false}; // NOLINT + std::atomic mConnectedEmitted{false}; // NOLINT std::recursive_mutex mMutex; // NOLINT // Only allow subclasses to be created @@ -124,7 +125,8 @@ class WebsocketConnection : public Connection, public EnableSharedFromThis { if (!self->mDestroyed) { if (self->mSocket && self->mSocket->readyState() == - rtc::WebSocket::State::Open) { + rtc::WebSocket::State::Open && + !self->mConnectedEmitted.exchange(true)) { SLogger::trace( "onOpen() emitting Connected " + getConnectionTypeString(), @@ -144,13 +146,25 @@ class WebsocketConnection : public Connection, public EnableSharedFromThis { SLogger::trace("setSocket() after move " + getConnectionTypeString()); - // Set socket callbacks + // Set socket callbacks. The message callback is deliberately NOT + // attached here — see startReceiving(). - mSocket->onMessage(this->onMessage, this->onStringMessage); mSocket->onError(this->onError); mSocket->onClosed(this->onClosed); mSocket->onOpen(this->onOpen); + // rtc does not retro-fire the open callback: if the websocket + // handshake completed on the processor thread before the callback + // above was attached, onOpen would never run and Connected would + // never be emitted. Emit it here in that case (mConnectedEmitted + // keeps the two paths idempotent). + if (mSocket->readyState() == rtc::WebSocket::State::Open) { + SLogger::trace( + "setSocket() socket already open, emitting Connected " + + getConnectionTypeString()); + this->onOpen(); + } + SLogger::trace("setSocket() end " + getConnectionTypeString()); } @@ -164,6 +178,24 @@ class WebsocketConnection : public Connection, public EnableSharedFromThis { destroy(); }; + // Attach the rtc message callback and deliver anything received so + // far. Deliberately separate from setSocket(): incoming frames queue + // inside libdatachannel until a message callback is attached (they + // flush synchronously at attach), while our Data event is + // fire-and-forget — a frame emitted before the application has + // registered its Data listener is silently lost. The owner therefore + // calls this only AFTER the application has had the opportunity to + // register listeners: the server after emitting Connected to the + // application, the client before open() (its listeners are + // registered before connect()). + void startReceiving() { + SLogger::trace("startReceiving() " + getConnectionTypeString()); + std::scoped_lock lock(mMutex); + if (!mDestroyed && mSocket) { + mSocket->onMessage(this->onMessage, this->onStringMessage); + } + } + void send(const std::vector& data) override { SLogger::trace( "send() start", @@ -195,20 +227,27 @@ class WebsocketConnection : public Connection, public EnableSharedFromThis { "close()", {{"connectionType", getConnectionTypeString()}, {"mDestroyed", mDestroyed.load()}}); - SLogger::debug( - "close() trying to acquire mutex lock in close()" + - getConnectionTypeString()); - std::scoped_lock lock(mMutex); - SLogger::debug( - "close() got mutex lock in close()" + getConnectionTypeString()); - if (!mDestroyed) { + std::shared_ptr socket; + { + std::scoped_lock lock(mMutex); + if (mDestroyed) { + SLogger::debug("close() on destroyed connection"); + return; + } mDestroyed = true; + socket = mSocket; + } + // The rtc calls happen OUTSIDE mMutex: rtc holds its callback + // mutex while a callback is executing, so resetCallbacks() blocks + // until any in-flight callback returns — and our rtc callbacks + // lock mMutex. Calling into rtc while holding mMutex is therefore + // a lock-order inversion (main thread: mMutex -> callback mutex; + // rtc thread: callback mutex -> mMutex) that deadlocked test + // teardowns. + if (socket) { SLogger::debug("close() resetting callbacks"); - mSocket->resetCallbacks(); - mSocket->close(); - // mSocket = nullptr; - } else { - SLogger::debug("close() on destroyed connection"); + socket->resetCallbacks(); + socket->close(); } SLogger::trace( "close() end", @@ -221,22 +260,20 @@ class WebsocketConnection : public Connection, public EnableSharedFromThis { "destroy()", {{"connectionType", getConnectionTypeString()}, {"mDestroyed", mDestroyed.load()}}); - SLogger::debug("destroy() trying to acquire mutex lock"); - std::scoped_lock lock(mMutex); - SLogger::debug("destroy() got mutex lock"); - if (mDestroyed) { - SLogger::debug("destroy() on destroyed connection"); - return; + std::shared_ptr socket; + { + std::scoped_lock lock(mMutex); + if (mDestroyed) { + SLogger::debug("destroy() on destroyed connection"); + return; + } + mDestroyed = true; + socket = mSocket; } - - mDestroyed = true; - if (mSocket) { - SLogger::debug("destroy() trying to get mutex lock"); - std::lock_guard lock(mMutex); - SLogger::debug("destroy() got mutex lock"); - mSocket->resetCallbacks(); - mSocket->close(); - // mSocket = nullptr; + // Outside mMutex for the same reason as in close(). + if (socket) { + socket->resetCallbacks(); + socket->close(); } } }; diff --git a/packages/streamr-dht/include/streamr-dht/connection/websocket/WebsocketServer.hpp b/packages/streamr-dht/include/streamr-dht/connection/websocket/WebsocketServer.hpp index ea9dbad6..c1d812c4 100644 --- a/packages/streamr-dht/include/streamr-dht/connection/websocket/WebsocketServer.hpp +++ b/packages/streamr-dht/include/streamr-dht/connection/websocket/WebsocketServer.hpp @@ -139,6 +139,14 @@ class WebsocketServer : public EventEmitter { readyConnection->removeAllListeners(); emit(readyConnection); + // Only now — after the application has registered its + // listeners on the connection inside the emit above — attach + // the rtc message callback. Frames the peer sent before this + // point have been queuing inside libdatachannel and flush + // here; attaching any earlier would emit Data with no + // listeners and silently drop the messages (the cause of the + // flaky Websocket/ConnectionLocking integration tests). + readyConnection->startReceiving(); SLogger::info("handleHalfReadySocket. Before erase"); mHalfReadyConnections.erase(id); });