From e6401c688b93ed3a688d13e32b6844dfcb22ce2f Mon Sep 17 00:00:00 2001 From: Ruslan Lesiutin Date: Mon, 3 Nov 2025 11:46:41 -0800 Subject: [PATCH] Limit WebSocket queue size for packager connection (#54300) Summary: # Changelog: [Internal] Establishes a queue mechanism on top of the OkHttp's WebSocket implementation. This mechanism will control the queue size and guarantee that we don't have more than 16MB scheduled. This prevents the scenario of when OkHttp forces WS disconnection because of this threshold. Reviewed By: motiz88 Differential Revision: D85581509 --- .../CxxInspectorPackagerConnection.kt | 108 ++++++++++++++++-- ...xxInspectorPackagerConnectionWebSocket.cpp | 31 ++++- .../CxxInspectorPackagerConnectionTest.kt | 28 +++++ 3 files changed, 156 insertions(+), 11 deletions(-) create mode 100644 packages/react-native/ReactAndroid/src/test/java/com/facebook/react/devsupport/CxxInspectorPackagerConnectionTest.kt diff --git a/packages/react-native/ReactAndroid/src/main/java/com/facebook/react/devsupport/CxxInspectorPackagerConnection.kt b/packages/react-native/ReactAndroid/src/main/java/com/facebook/react/devsupport/CxxInspectorPackagerConnection.kt index a92d4ac26c62..5e3be8b14413 100644 --- a/packages/react-native/ReactAndroid/src/main/java/com/facebook/react/devsupport/CxxInspectorPackagerConnection.kt +++ b/packages/react-native/ReactAndroid/src/main/java/com/facebook/react/devsupport/CxxInspectorPackagerConnection.kt @@ -9,11 +9,17 @@ package com.facebook.react.devsupport import android.os.Handler import android.os.Looper +import com.facebook.common.logging.FLog import com.facebook.jni.HybridData import com.facebook.proguard.annotations.DoNotStrip import com.facebook.proguard.annotations.DoNotStripAny +import com.facebook.react.common.annotations.VisibleForTesting import com.facebook.soloader.SoLoader import java.io.Closeable +import java.nio.ByteBuffer +import java.nio.charset.StandardCharsets +import java.util.ArrayDeque +import java.util.Queue import java.util.concurrent.TimeUnit import okhttp3.OkHttpClient import okhttp3.Request @@ -67,7 +73,7 @@ internal class CxxInspectorPackagerConnection( */ @DoNotStripAny private interface IWebSocket : Closeable { - fun send(message: String) + fun send(chunk: ByteBuffer) /** * Close the WebSocket connection. NOTE: There is no close() method in the C++ interface. @@ -76,6 +82,95 @@ internal class CxxInspectorPackagerConnection( override fun close() } + /** + * A simple WebSocket wrapper that prevents having more than 16MiB of messages queued + * simultaneously. This is done to stop OkHttp from closing the WebSocket connection. + * + * https://github.com/facebook/react-native/issues/39651. + * https://github.com/square/okhttp/blob/4e7dbec1ea6c9cf8d80422ac9d44b9b185c749a3/okhttp/src/commonJvmAndroid/kotlin/okhttp3/internal/ws/RealWebSocket.kt#L684. + */ + private class InspectorPackagerWebSocketImpl( + private val nativeWebSocket: WebSocket, + private val handler: Handler, + ) : IWebSocket { + private val messageQueue: Queue> = ArrayDeque() + private val queueLock = Any() + private val drainRunnable = + object : Runnable { + override fun run() { + FLog.d(TAG, "Attempting to drain the message queue after ${drainDelayMs}ms") + tryDrainQueue() + } + } + + /** + * We are providing a String to OkHttp's WebSocket, because there is no guarantee that all CDP + * clients will support binary data format. + */ + override fun send(chunk: ByteBuffer) { + synchronized(queueLock) { + val messageSize = chunk.capacity() + val message = StandardCharsets.UTF_8.decode(chunk).toString() + val currentQueueSize = nativeWebSocket.queueSize() + + if (currentQueueSize + messageSize > MAX_QUEUE_SIZE) { + FLog.d(TAG, "Reached queue size limit. Queueing the message.") + messageQueue.offer(Pair(message, messageSize)) + scheduleDrain() + } else { + if (messageQueue.isEmpty()) { + nativeWebSocket.send(message) + } else { + messageQueue.offer(Pair(message, messageSize)) + tryDrainQueue() + } + } + } + } + + override fun close() { + synchronized(queueLock) { + handler.removeCallbacks(drainRunnable) + messageQueue.clear() + nativeWebSocket.close(1000, "End of session") + } + } + + private fun tryDrainQueue() { + synchronized(queueLock) { + while (messageQueue.isNotEmpty()) { + val (nextMessage, nextMessageSize) = messageQueue.peek() ?: break + val currentQueueSize = nativeWebSocket.queueSize() + + if (currentQueueSize + nextMessageSize <= MAX_QUEUE_SIZE) { + messageQueue.poll() + if (!nativeWebSocket.send(nextMessage)) { + // The WebSocket is closing, closed, or cancelled. + handler.removeCallbacks(drainRunnable) + messageQueue.clear() + + break + } + } else { + scheduleDrain() + break + } + } + } + } + + private fun scheduleDrain() { + FLog.d(TAG, "Scheduled a task to drain messages queue.") + handler.removeCallbacks(drainRunnable) + handler.postDelayed(drainRunnable, drainDelayMs) + } + + companion object { + private val TAG: String = InspectorPackagerWebSocketImpl::class.java.simpleName + private const val drainDelayMs: Long = 100 + } + } + /** Java implementation of the C++ InspectorPackagerConnectionDelegate interface. */ private class DelegateImpl { private val httpClient = @@ -130,15 +225,8 @@ internal class CxxInspectorPackagerConnection( } }, ) - return object : IWebSocket { - override fun send(message: String) { - webSocket.send(message) - } - override fun close() { - webSocket.close(1000, "End of session") - } - } + return InspectorPackagerWebSocketImpl(webSocket, handler) } @DoNotStrip @@ -152,6 +240,8 @@ internal class CxxInspectorPackagerConnection( SoLoader.loadLibrary("react_devsupportjni") } + @VisibleForTesting internal const val MAX_QUEUE_SIZE = 16L * 1024 * 1024 // 16MiB + @JvmStatic private external fun initHybrid( url: String, diff --git a/packages/react-native/ReactAndroid/src/main/jni/react/devsupport/JCxxInspectorPackagerConnectionWebSocket.cpp b/packages/react-native/ReactAndroid/src/main/jni/react/devsupport/JCxxInspectorPackagerConnectionWebSocket.cpp index 7de472736303..efd1948f7aa5 100644 --- a/packages/react-native/ReactAndroid/src/main/jni/react/devsupport/JCxxInspectorPackagerConnectionWebSocket.cpp +++ b/packages/react-native/ReactAndroid/src/main/jni/react/devsupport/JCxxInspectorPackagerConnectionWebSocket.cpp @@ -5,6 +5,8 @@ * LICENSE file in the root directory of this source tree. */ +#include + #include "JCxxInspectorPackagerConnectionWebSocket.h" using namespace facebook::jni; @@ -12,10 +14,35 @@ using namespace facebook::react::jsinspector_modern; namespace facebook::react::jsinspector_modern { +namespace { + +local_ref getReadOnlyByteBufferFromStringView( + std::string_view sv) { + auto buffer = JByteBuffer::wrapBytes( + const_cast(reinterpret_cast(sv.data())), + sv.size()); + + /** + * Return a read-only buffer that shares the underlying contents. + * This guards from accidential mutations on the Java side, since we did + * casting above. + * + * https://docs.oracle.com/javase/8/docs/api/java/nio/ByteBuffer.html#asReadOnlyBuffer-- + */ + static auto method = + buffer->javaClassStatic()->getMethod( + "asReadOnlyBuffer"); + return method(buffer); +} + +} // namespace + void JCxxInspectorPackagerConnectionWebSocket::send(std::string_view message) { static auto method = - javaClassStatic()->getMethod("send"); - method(self(), std::string(message)); + javaClassStatic()->getMethod)>( + "send"); + auto byteBuffer = getReadOnlyByteBufferFromStringView(message); + method(self(), byteBuffer); } void JCxxInspectorPackagerConnectionWebSocket::close() { diff --git a/packages/react-native/ReactAndroid/src/test/java/com/facebook/react/devsupport/CxxInspectorPackagerConnectionTest.kt b/packages/react-native/ReactAndroid/src/test/java/com/facebook/react/devsupport/CxxInspectorPackagerConnectionTest.kt new file mode 100644 index 000000000000..462317018ec7 --- /dev/null +++ b/packages/react-native/ReactAndroid/src/test/java/com/facebook/react/devsupport/CxxInspectorPackagerConnectionTest.kt @@ -0,0 +1,28 @@ +/* + * Copyright (c) Meta Platforms, Inc. and affiliates. + * + * This source code is licensed under the MIT license found in the + * LICENSE file in the root directory of this source tree. + */ + +package com.facebook.react.devsupport + +import okhttp3.internal.ws.RealWebSocket +import org.assertj.core.api.Assertions.assertThat +import org.junit.Test + +class CxxInspectorPackagerConnectionTest { + + @Test + fun testMaxQueueSizeEquality() { + val okHttpRealWebSocketClass = RealWebSocket::class.java + val okHttpMaxQueueSizeField = okHttpRealWebSocketClass.getDeclaredField("MAX_QUEUE_SIZE") + okHttpMaxQueueSizeField.isAccessible = true + + val okHttpMaxQueueSize = okHttpMaxQueueSizeField.getLong(null) + assertThat(okHttpMaxQueueSize).isNotNull + + assertThat(okHttpMaxQueueSize) + .isEqualTo(CxxInspectorPackagerConnection.Companion.MAX_QUEUE_SIZE) + } +}