From 416beab5d8cb5d6806219cc33b2231d1531ea4dc Mon Sep 17 00:00:00 2001 From: Andre Dietisheim Date: Thu, 24 Sep 2026 17:07:19 +0200 Subject: [PATCH] enhancement: harden Java port-forward path (CRW-13358) * Long-lived PF accept/copy on Dispatchers.IO starves thread pool and amplifies remote UI lag. Use dedicated threads instead. * TCP_NODELAY, 64KiB copy without WebSocket, per-chunk flush reduce lag. Signed-off-by: Andre Dietisheim Co-authored-by: Cursor --- README.md | 8 ++++ .../gateway/devworkspace/DevWorkspaces.kt | 2 +- .../gateway/openshift/DevWorkspacePods.kt | 46 +++++++++++++------ .../openshift/PortForwardDispatcher.kt | 39 ++++++++++++++++ 4 files changed, 81 insertions(+), 14 deletions(-) create mode 100644 src/main/kotlin/com/redhat/devtools/gateway/openshift/PortForwardDispatcher.kt diff --git a/README.md b/README.md index 2139302a..1d696079 100644 --- a/README.md +++ b/README.md @@ -66,6 +66,14 @@ sudo lsof -i -P | grep LISTEN | grep 5990 2. The Gateway application and the Dev Spaces plugin logs are stored in `/Users//Library/Logs/JetBrains/JetBrainsGateway/idea.log` +3. Port-forward transport (Java WebSocket): + - On connect, logs should include: + `Starting port forward on local port … (transport=java-websocket, copyBuffer=65536, dispatcher=devspaces-port-forward)` + - If the UI feels sluggish or you see “IntelliJ IDEA has encountered a slowdown”, collect: + - Gateway/IDEA logs (`…/Library/Logs/JetBrains/JetBrainsGateway*/idea.log` or IDEA equivalent) + - Any `PerformanceWatcherImpl` / `Dispatchers.IO` / thread-dump lines + - Whether the banner appears during idle editing vs only while `Analyzing...` + ## Release - Find a draft release on the [Releases](https://github.com/redhat-developer/devspaces-gateway-plugin/releases) page. The draft is created and updated automatically on each push to the `main` branch. diff --git a/src/main/kotlin/com/redhat/devtools/gateway/devworkspace/DevWorkspaces.kt b/src/main/kotlin/com/redhat/devtools/gateway/devworkspace/DevWorkspaces.kt index 233802c3..3b6faf8c 100644 --- a/src/main/kotlin/com/redhat/devtools/gateway/devworkspace/DevWorkspaces.kt +++ b/src/main/kotlin/com/redhat/devtools/gateway/devworkspace/DevWorkspaces.kt @@ -158,7 +158,7 @@ class DevWorkspaces(private val client: ApiClient) { } } -@Throws(ApiException::class) + @Throws(ApiException::class) fun start(namespace: String, name: String) { DevWorkspacePatch(namespace, name, client) { get(namespace, name) diff --git a/src/main/kotlin/com/redhat/devtools/gateway/openshift/DevWorkspacePods.kt b/src/main/kotlin/com/redhat/devtools/gateway/openshift/DevWorkspacePods.kt index c1e83a8a..d1a2f6cb 100644 --- a/src/main/kotlin/com/redhat/devtools/gateway/openshift/DevWorkspacePods.kt +++ b/src/main/kotlin/com/redhat/devtools/gateway/openshift/DevWorkspacePods.kt @@ -1,5 +1,5 @@ /* - * Copyright (c) 2024-2025 Red Hat, Inc. + * Copyright (c) 2024-2026 Red Hat, Inc. * This program and the accompanying materials are made * available under the terms of the Eclipse Public License 2.0 * which is available at https://www.eclipse.org/legal/epl-2.0/ @@ -30,6 +30,7 @@ class DevWorkspacePods(private val client: ApiClient) { companion object { const val WORKSPACE_LABEL_KEY = "controller.devfile.io/devworkspace_name" + private const val COPY_BUFFER_SIZE = 64 * 1024 private const val CONNECT_ATTEMPTS = 5 private const val RECONNECT_DELAY: Long = 1000 } @@ -73,8 +74,8 @@ class DevWorkspacePods(private val client: ApiClient) { fun forward(pod: V1Pod, localPort: Int, remotePort: Int): Closeable { val serverSocket = ServerSocket(localPort, 50, InetAddress.getLoopbackAddress()) val scope = CoroutineScope( - // dont cancel if child coroutine fails + use blocking I/O scope - SupervisorJob() + Dispatchers.IO + // dont cancel if child coroutine fails + use dedicated port-forward dispatcher + SupervisorJob() + PortForwardDispatcher.dispatcher ) scope.acceptConnections(serverSocket, pod, localPort, remotePort) return Closeable { @@ -90,10 +91,14 @@ class DevWorkspacePods(private val client: ApiClient) { remotePort: Int ) { launch { - logger.info("Starting port forward on local port $localPort...") + logger.info( + "Starting port forward on local port $localPort " + + "(transport=java-websocket, copyBuffer=$COPY_BUFFER_SIZE, dispatcher=devspaces-port-forward)" + ) while (isActive) { val clientSocket = createClientSocket(serverSocket) ?: break + applyClientSocketOptions(clientSocket) launch { handleConnection( @@ -170,7 +175,7 @@ class DevWorkspacePods(private val client: ApiClient) { ensureActive() launch { try { - clientSocket.getInputStream().copyToAndFlush(forwardResult.getOutboundStream(remotePort)) + clientSocket.getInputStream().copyToWebSocketOutbound(forwardResult.getOutboundStream(remotePort)) } catch (e: Exception) { closeStreams(remotePort, forwardResult) throw e @@ -178,7 +183,7 @@ class DevWorkspacePods(private val client: ApiClient) { } launch { try { - forwardResult.getInputStream(remotePort).copyToAndFlush(clientSocket.getOutputStream()) + forwardResult.getInputStream(remotePort).copyToLocalSocket(clientSocket.getOutputStream()) } catch (e: Exception) { closeStreams(remotePort, forwardResult) throw e @@ -197,14 +202,29 @@ class DevWorkspacePods(private val client: ApiClient) { .onFailure { logger.debug("Could not get outbound stream for port $port while closing port-forward", it) } } - private fun InputStream.copyToAndFlush(destination: OutputStream) { - try { - copyTo(destination) - destination.flush() - } catch (e: IOException) { - logger.info("IOException during stream copy.", e) - throw e + private fun applyClientSocketOptions(socket: Socket) { + socket.tcpNoDelay = true + } + + private fun InputStream.copyToLocalSocket(destination: OutputStream) { + val buffer = ByteArray(COPY_BUFFER_SIZE) + while (true) { + val n = read(buffer) + if (n < 0) break + destination.write(buffer, 0, n) + destination.flush() // local socket ONLY + } + } + + private fun InputStream.copyToWebSocketOutbound(destination: OutputStream) { + val buffer = ByteArray(COPY_BUFFER_SIZE) + while (true) { + val n = read(buffer) + if (n < 0) break + destination.write(buffer, 0, n) + // do NOT flush per chunk — client-java 24 WebSocketOutputStream.flush() sleeps 100ms } + runCatching { destination.flush() } } @Throws(IOException::class) diff --git a/src/main/kotlin/com/redhat/devtools/gateway/openshift/PortForwardDispatcher.kt b/src/main/kotlin/com/redhat/devtools/gateway/openshift/PortForwardDispatcher.kt new file mode 100644 index 00000000..3fea3b1c --- /dev/null +++ b/src/main/kotlin/com/redhat/devtools/gateway/openshift/PortForwardDispatcher.kt @@ -0,0 +1,39 @@ +/* + * Copyright (c) 2026 Red Hat, Inc. + * This program and the accompanying materials are made + * available under the terms of the Eclipse Public License 2.0 + * which is available at https://www.eclipse.org/legal/epl-2.0/ + * + * SPDX-License-Identifier: EPL-2.0 + * + * Contributors: + * Red Hat, Inc. - initial API and implementation + */ +package com.redhat.devtools.gateway.openshift + +import kotlinx.coroutines.CoroutineDispatcher +import kotlinx.coroutines.asCoroutineDispatcher +import java.util.concurrent.Executors +import java.util.concurrent.atomic.AtomicInteger + +/** + * Dedicated dispatcher for long-lived port-forward accept/copy loops. + * + * Must not use `Dispatchers.IO`: port-forward sessions block indefinitely on + * accept/copy, and running them on the shared IO pool would starve other IO + * work (CRW-12992). Backed by a cached thread pool of daemon threads named + * `devspaces-port-forward-${n}` so long-lived forward sessions can be + * identified in thread dumps. + */ +internal object PortForwardDispatcher { + + private val threadCounter = AtomicInteger(0) + + val dispatcher: CoroutineDispatcher = Executors + .newCachedThreadPool { runnable -> + Thread(runnable, "devspaces-port-forward-${threadCounter.incrementAndGet()}").apply { + isDaemon = true + } + } + .asCoroutineDispatcher() +}