From f1c980001aaece74449466a181c52ac443bbd5e1 Mon Sep 17 00:00:00 2001 From: Roman Herbstmann Date: Thu, 10 Sep 2026 23:16:36 +0200 Subject: [PATCH 1/2] rtmp/rtsp: unblock a sender stuck in a socket write on disconnect With TCP backpressure (server ACKs not arriving, send window full) the sender blocks in a java.io socket write, which ignores coroutine cancellation. disconnect() stopped the sender first (BaseSender.stop -> job.cancelAndJoin) and closed the socket only afterwards, so the join waited until the network recovered or TCP gave up. A reConnect hung for the whole time. Try a cooperative stop first (keeps the graceful close in the normal case) and, if the sender does not stop within 1 s, close the socket to unblock the write and stop again. RTSP over TCP (interleaved) writes to the same socket and has the same issue. SRT/UDP are not affected (UDP sends do not block on ACKs). --- .../main/java/com/pedro/rtmp/rtmp/RtmpClient.kt | 14 +++++++++++++- .../main/java/com/pedro/rtsp/rtsp/RtspClient.kt | 14 +++++++++++++- 2 files changed, 26 insertions(+), 2 deletions(-) diff --git a/rtmp/src/main/java/com/pedro/rtmp/rtmp/RtmpClient.kt b/rtmp/src/main/java/com/pedro/rtmp/rtmp/RtmpClient.kt index 30d11b879..2de66e272 100644 --- a/rtmp/src/main/java/com/pedro/rtmp/rtmp/RtmpClient.kt +++ b/rtmp/src/main/java/com/pedro/rtmp/rtmp/RtmpClient.kt @@ -59,6 +59,9 @@ import java.nio.ByteBuffer import java.util.concurrent.atomic.AtomicLong import javax.net.ssl.TrustManager import kotlin.time.Duration.Companion.milliseconds +import kotlin.time.Duration.Companion.seconds + +private val SENDER_STOP_TIMEOUT = 1.seconds /** * Created by pedro on 8/04/21. @@ -566,7 +569,16 @@ class RtmpClient(private val connectChecker: ConnectChecker) { } private suspend fun disconnect(clear: Boolean) { - if (isStreaming) rtmpSender.stop(clear) + if (isStreaming) { + //the sender can be blocked in a socket write (TCP backpressure) that ignores cancellation, + //only closing the socket unblocks it. Try a cooperative stop first to keep the graceful close. + val stopped = withTimeoutOrNull(SENDER_STOP_TIMEOUT) { rtmpSender.stop(clear) } != null + if (!stopped) { + Log.w(TAG, "sender blocked in socket write, closing socket to unblock it") + runCatching { socket?.close() } + rtmpSender.stop(clear) + } + } runCatching { withTimeoutOrNull(100.milliseconds) { socket?.let { commandsManager.sendClose(it) } diff --git a/rtsp/src/main/java/com/pedro/rtsp/rtsp/RtspClient.kt b/rtsp/src/main/java/com/pedro/rtsp/rtsp/RtspClient.kt index 67720bddb..0ef4a8f5e 100644 --- a/rtsp/src/main/java/com/pedro/rtsp/rtsp/RtspClient.kt +++ b/rtsp/src/main/java/com/pedro/rtsp/rtsp/RtspClient.kt @@ -46,6 +46,9 @@ import java.net.URISyntaxException import java.nio.ByteBuffer import javax.net.ssl.TrustManager import kotlin.time.Duration.Companion.milliseconds +import kotlin.time.Duration.Companion.seconds + +private val SENDER_STOP_TIMEOUT = 1.seconds /** * Created by pedro on 10/02/17. @@ -426,7 +429,16 @@ class RtspClient(private val connectChecker: ConnectChecker) { } private suspend fun disconnect(clear: Boolean) { - if (isStreaming) rtspSender.stop() + if (isStreaming) { + //the sender can be blocked in a socket write (TCP backpressure) that ignores cancellation, + //only closing the socket unblocks it. Try a cooperative stop first to keep the graceful close. + val stopped = withTimeoutOrNull(SENDER_STOP_TIMEOUT) { rtspSender.stop() } != null + if (!stopped) { + Log.w(TAG, "sender blocked in socket write, closing socket to unblock it") + runCatching { socket?.close() } + rtspSender.stop() + } + } val error = runCatching { withTimeoutOrNull(100.milliseconds) { socket?.write(commandsManager.createTeardown()) From cf419bd21ecd7f4fd2b5c7bf7d5d6d3b23dfbe87 Mon Sep 17 00:00:00 2001 From: pedroSG94 Date: Fri, 11 Sep 2026 21:24:19 +0200 Subject: [PATCH 2/2] refactor sender stop to delegate close into stop sender --- .../java/com/pedro/common/base/BaseSender.kt | 10 ++++++++-- .../java/com/pedro/rtmp/rtmp/RtmpClient.kt | 16 ++------------- .../java/com/pedro/rtsp/rtsp/RtspClient.kt | 20 ++++--------------- 3 files changed, 14 insertions(+), 32 deletions(-) diff --git a/common/src/main/java/com/pedro/common/base/BaseSender.kt b/common/src/main/java/com/pedro/common/base/BaseSender.kt index 7aa44e7e2..1c4cb8be6 100644 --- a/common/src/main/java/com/pedro/common/base/BaseSender.kt +++ b/common/src/main/java/com/pedro/common/base/BaseSender.kt @@ -18,8 +18,10 @@ import kotlinx.coroutines.delay import kotlinx.coroutines.isActive import kotlinx.coroutines.launch import kotlinx.coroutines.runInterruptible +import kotlinx.coroutines.withTimeoutOrNull import java.nio.ByteBuffer import java.util.concurrent.atomic.AtomicLong +import kotlin.time.Duration.Companion.milliseconds abstract class BaseSender( protected val connectChecker: ConnectChecker, @@ -119,7 +121,7 @@ abstract class BaseSender( } } - suspend fun stop(clear: Boolean = true) { + suspend fun stop(clear: Boolean = true, unlockNeeded: suspend () -> Unit = {}) { running = false stopImp(clear) resetSentAudioFrames() @@ -127,7 +129,11 @@ abstract class BaseSender( resetDroppedAudioFrames() resetDroppedVideoFrames() resetBytesSend() - job?.cancelAndJoin() + val stopped = withTimeoutOrNull(1000.milliseconds) { job?.cancelAndJoin() } != null + if (!stopped) { + unlockNeeded() + withTimeoutOrNull(1000.milliseconds) { job?.cancelAndJoin() } + } job = null queue.clear { bufferPool.release(it.data) } bufferPool.clear() diff --git a/rtmp/src/main/java/com/pedro/rtmp/rtmp/RtmpClient.kt b/rtmp/src/main/java/com/pedro/rtmp/rtmp/RtmpClient.kt index 2de66e272..43f3d30ce 100644 --- a/rtmp/src/main/java/com/pedro/rtmp/rtmp/RtmpClient.kt +++ b/rtmp/src/main/java/com/pedro/rtmp/rtmp/RtmpClient.kt @@ -33,11 +33,11 @@ import com.pedro.common.validMessage import com.pedro.rtmp.rtmp.message.Abort import com.pedro.rtmp.rtmp.message.Acknowledgement import com.pedro.rtmp.rtmp.message.Aggregate +import com.pedro.rtmp.rtmp.message.Command import com.pedro.rtmp.rtmp.message.MessageType import com.pedro.rtmp.rtmp.message.SetChunkSize import com.pedro.rtmp.rtmp.message.SetPeerBandwidth import com.pedro.rtmp.rtmp.message.WindowAcknowledgementSize -import com.pedro.rtmp.rtmp.message.Command import com.pedro.rtmp.rtmp.message.control.Type import com.pedro.rtmp.rtmp.message.control.UserControl import com.pedro.rtmp.utils.AuthUtil @@ -59,9 +59,6 @@ import java.nio.ByteBuffer import java.util.concurrent.atomic.AtomicLong import javax.net.ssl.TrustManager import kotlin.time.Duration.Companion.milliseconds -import kotlin.time.Duration.Companion.seconds - -private val SENDER_STOP_TIMEOUT = 1.seconds /** * Created by pedro on 8/04/21. @@ -569,16 +566,7 @@ class RtmpClient(private val connectChecker: ConnectChecker) { } private suspend fun disconnect(clear: Boolean) { - if (isStreaming) { - //the sender can be blocked in a socket write (TCP backpressure) that ignores cancellation, - //only closing the socket unblocks it. Try a cooperative stop first to keep the graceful close. - val stopped = withTimeoutOrNull(SENDER_STOP_TIMEOUT) { rtmpSender.stop(clear) } != null - if (!stopped) { - Log.w(TAG, "sender blocked in socket write, closing socket to unblock it") - runCatching { socket?.close() } - rtmpSender.stop(clear) - } - } + if (isStreaming) rtmpSender.stop(clear, unlockNeeded = { socket?.close() }) runCatching { withTimeoutOrNull(100.milliseconds) { socket?.let { commandsManager.sendClose(it) } diff --git a/rtsp/src/main/java/com/pedro/rtsp/rtsp/RtspClient.kt b/rtsp/src/main/java/com/pedro/rtsp/rtsp/RtspClient.kt index 0ef4a8f5e..1ec0ae001 100644 --- a/rtsp/src/main/java/com/pedro/rtsp/rtsp/RtspClient.kt +++ b/rtsp/src/main/java/com/pedro/rtsp/rtsp/RtspClient.kt @@ -46,9 +46,6 @@ import java.net.URISyntaxException import java.nio.ByteBuffer import javax.net.ssl.TrustManager import kotlin.time.Duration.Companion.milliseconds -import kotlin.time.Duration.Companion.seconds - -private val SENDER_STOP_TIMEOUT = 1.seconds /** * Created by pedro on 10/02/17. @@ -429,25 +426,16 @@ class RtspClient(private val connectChecker: ConnectChecker) { } private suspend fun disconnect(clear: Boolean) { - if (isStreaming) { - //the sender can be blocked in a socket write (TCP backpressure) that ignores cancellation, - //only closing the socket unblocks it. Try a cooperative stop first to keep the graceful close. - val stopped = withTimeoutOrNull(SENDER_STOP_TIMEOUT) { rtspSender.stop() } != null - if (!stopped) { - Log.w(TAG, "sender blocked in socket write, closing socket to unblock it") - runCatching { socket?.close() } - rtspSender.stop() - } - } + if (isStreaming) rtspSender.stop(unlockNeeded = { socket?.close() }) val error = runCatching { withTimeoutOrNull(100.milliseconds) { socket?.write(commandsManager.createTeardown()) socket?.flush() + Log.i(TAG, "write teardown success") } - socket?.close() - socket = null - Log.i(TAG, "write teardown success") }.exceptionOrNull() + socket?.close() + socket = null if (error != null) { Log.e(TAG, "disconnect error", error) }