From 429adc01736fcbb8cc11e0bc22f26b70f7467020 Mon Sep 17 00:00:00 2001 From: ThibaultBee <37510686+ThibaultBee@users.noreply.github.com> Date: Sat, 26 Sep 2026 21:39:08 +0200 Subject: [PATCH] feat(rtmp): add ping functionality and round trip time measurement to RtmpRawMetrics --- .../rtmp/elements/endpoints/RtmpEndpoint.kt | 7 ++- .../elements/endpoints/RtmpEndpointMetrics.kt | 52 +++++++++++++++---- gradle/libs.versions.toml | 2 +- 3 files changed, 50 insertions(+), 11 deletions(-) diff --git a/extensions/rtmp/src/main/java/io/github/thibaultbee/streampack/ext/rtmp/elements/endpoints/RtmpEndpoint.kt b/extensions/rtmp/src/main/java/io/github/thibaultbee/streampack/ext/rtmp/elements/endpoints/RtmpEndpoint.kt index 302e9c53d..26005e4d7 100644 --- a/extensions/rtmp/src/main/java/io/github/thibaultbee/streampack/ext/rtmp/elements/endpoints/RtmpEndpoint.kt +++ b/extensions/rtmp/src/main/java/io/github/thibaultbee/streampack/ext/rtmp/elements/endpoints/RtmpEndpoint.kt @@ -123,8 +123,13 @@ class RtmpEndpoint internal constructor( } } + private val rtmpRawMetrics = RtmpRawMetrics( + clientProvider = { rtmpClient }, + metricsProvider = { syncMetrics } + ) + override val metrics: RtmpEndpointMetrics - get() = RtmpEndpointMetrics { syncMetrics } + get() = RtmpEndpointMetrics(rtmpRawMetrics) private val _isOpenFlow = MutableStateFlow(false) override val isOpenFlow = _isOpenFlow.asStateFlow() diff --git a/extensions/rtmp/src/main/java/io/github/thibaultbee/streampack/ext/rtmp/elements/endpoints/RtmpEndpointMetrics.kt b/extensions/rtmp/src/main/java/io/github/thibaultbee/streampack/ext/rtmp/elements/endpoints/RtmpEndpointMetrics.kt index 93609a04b..154c99844 100644 --- a/extensions/rtmp/src/main/java/io/github/thibaultbee/streampack/ext/rtmp/elements/endpoints/RtmpEndpointMetrics.kt +++ b/extensions/rtmp/src/main/java/io/github/thibaultbee/streampack/ext/rtmp/elements/endpoints/RtmpEndpointMetrics.kt @@ -15,16 +15,14 @@ */ package io.github.thibaultbee.streampack.ext.rtmp.elements.endpoints +import io.github.komedia.komuxer.rtmp.client.RtmpClient +import io.github.komedia.komuxer.rtmp.messages.UserControl import io.github.komedia.komuxer.rtmp.util.metrics.RtmpMetrics import io.github.thibaultbee.streampack.core.elements.metrics.EndpointMetrics +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock import kotlin.time.Duration - -/** - * Creates a [RtmpEndpointMetrics] from a [metricsProvider]. - */ -fun RtmpEndpointMetrics(metricsProvider: () -> RtmpMetrics?): RtmpEndpointMetrics { - return RtmpEndpointMetrics(RtmpRawMetrics(metricsProvider)) -} +import kotlin.time.measureTime /** * Creates a [RtmpEndpointMetrics] from a [RtmpMetrics]. @@ -59,12 +57,48 @@ data class RtmpEndpointMetrics( /** * Provides an access to internal RTMP metrics APIs. */ -class RtmpRawMetrics internal constructor(private val metricsProvider: () -> RtmpMetrics?) { +class RtmpRawMetrics internal constructor( + private val clientProvider: () -> RtmpClient?, + private val metricsProvider: () -> RtmpMetrics? +) { + private val pingMutex = Mutex() + /** * Returns the [RtmpMetrics] if the client is available, otherwise null. */ val rtmpMetricsOrNull: RtmpMetrics? - get() = metricsProvider() + get() = metricsProvider() ?: clientProvider()?.metrics + + private suspend fun writePingInternal(): UserControl { + val client = clientProvider() ?: throw IllegalStateException("RTMP client is not available") + if (client.isClosed) { + throw IllegalStateException("RTMP client is closed") + } + return client.writePing() + } + + /** + * Writes a ping request to the server and awaits the response. + * + * @return the ping response [UserControl] + * @throws IllegalStateException if the RTMP client is not available or closed + */ + suspend fun writePing(): UserControl = pingMutex.withLock { + writePingInternal() + } + + /** + * Computes the round trip time (RTT) to the RTMP server using a ping request. + * + * If the server has not implemented the ping response, it will hang indefinitely. + * + * @return the measured [Duration], or null if the client is not available, closed, or an error occurred. + */ + suspend fun rtt(): Duration = pingMutex.withLock { + measureTime { + writePingInternal() + } + } } /** diff --git a/gradle/libs.versions.toml b/gradle/libs.versions.toml index f33814572..bf2684323 100644 --- a/gradle/libs.versions.toml +++ b/gradle/libs.versions.toml @@ -26,7 +26,7 @@ kotlinxIo = "0.9.1" material = "1.14.0" mockk = "1.14.11" robolectric = "4.16.1" -komuxer = "0.4.0" +komuxer = "0.4.2" srtdroid = "1.10.1" junitKtx = "1.3.0" compose = "1.12.0"