From 624431a0b8312e5e82657277c879cadac319eadc Mon Sep 17 00:00:00 2001 From: rahullohra Date: Thu, 6 Aug 2026 16:29:51 +0530 Subject: [PATCH 1/2] fix: Fix incorrect reporting of was pc ever connected --- .../call/observer/PeerConnectionAnalytics.kt | 12 ++++ .../PeerConnectionAnalyticsStateHolder.kt | 19 ++++++ .../reporting/ClientEventReporter.kt | 7 +- .../observer/PeerConnectionAnalyticsTest.kt | 65 +++++++++++++++++++ .../reporting/ClientEventReporterTest.kt | 3 + 5 files changed, 102 insertions(+), 4 deletions(-) diff --git a/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalytics.kt b/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalytics.kt index e01dfe68cef..e6f064cc290 100644 --- a/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalytics.kt +++ b/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalytics.kt @@ -162,6 +162,17 @@ internal class PeerConnectionAnalytics( iceState: VideoAnalyticsIceState, peerConnectionState: PeerConnection.PeerConnectionState?, ) { + when (peerConnectionState) { + PeerConnection.PeerConnectionState.CONNECTED -> { + if (role == PeerConnectionRole.PUBLISH) { + stateHolder.updatePublisherEverConnected(true) + } else { + stateHolder.updateSubscriberEverConnected(true) + } + } + else -> {} + } + reporter.onPeerConnectionStateChanged( peerConnectionHashCode = peerConnectionHashCode, callId = callId, @@ -172,6 +183,7 @@ internal class PeerConnectionAnalytics( joinStageAttemptId = joinAnalyticsStateHolder.state.value.joinStageAttemptId ?: "unknown", joinReason = joinAnalyticsStateHolder.state.value.joinReason ?: JoinReason.Unknown, sfuId = sfuAnalyticsStateHolder.sfuId.value, + wasPrevConnected = stateHolder.isPcEverConnected(role), callSessionId = joinAnalyticsStateHolder.state.value.callSessionId, ) } diff --git a/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalyticsStateHolder.kt b/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalyticsStateHolder.kt index e16ac1cce89..b51e392007a 100644 --- a/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalyticsStateHolder.kt +++ b/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalyticsStateHolder.kt @@ -17,6 +17,7 @@ package io.getstream.video.android.core.analytics.call.observer import io.getstream.video.android.core.analytics.call.observer.model.Stage +import io.getstream.video.android.core.analytics.reporting.model.PeerConnectionRole import kotlinx.coroutines.Job import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.StateFlow @@ -47,6 +48,22 @@ internal class PeerConnectionAnalyticsStateHolder { _state.update { it.copy(subscriberStage = stage) } } + fun updatePublisherEverConnected(wasPublisherEverConnected: Boolean) { + _state.update { it.copy(wasPublisherEverConnected = wasPublisherEverConnected) } + } + + fun updateSubscriberEverConnected(wasSubscriberEverConnected: Boolean) { + _state.update { it.copy(wasSubscriberEverConnected = wasSubscriberEverConnected) } + } + + fun isPcEverConnected(role: PeerConnectionRole): Boolean { + return if (role == PeerConnectionRole.PUBLISH) { + _state.value.wasPublisherEverConnected + } else { + _state.value.wasSubscriberEverConnected + } + } + fun update( peerConnectionObserverJob: Job? = state.value.peerConnectionObserverJob, publisherJob: Job? = state.value.publisherJob, @@ -76,4 +93,6 @@ internal data class PeerConnectionAnalyticsState( val subscriberJob: Job? = null, val publisherStage: Stage = Stage.NOT_STARTED, val subscriberStage: Stage = Stage.NOT_STARTED, + val wasPublisherEverConnected: Boolean = false, + val wasSubscriberEverConnected: Boolean = false, ) diff --git a/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/analytics/reporting/ClientEventReporter.kt b/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/analytics/reporting/ClientEventReporter.kt index e362dd5e627..11ecf7b62d4 100644 --- a/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/analytics/reporting/ClientEventReporter.kt +++ b/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/analytics/reporting/ClientEventReporter.kt @@ -82,7 +82,6 @@ internal class ClientEventReporter( ClientEventFactory(sdkVersion, userAgent, coordinatorAnalyticsStateHolder) private val postCallFlightSessions = ConcurrentHashMap() - private val pcEverConnected = ConcurrentHashMap() private val pcEventReporterStateHolder = PeerConnectionEventReporterStateHolder() // --- Coordinator WS --- @@ -285,6 +284,7 @@ internal class ClientEventReporter( joinReason: JoinReason, role: PeerConnectionRole, iceState: VideoAnalyticsIceState, + wasPrevConnected: Boolean, peerConnectionState: PeerConnection.PeerConnectionState?, ) { when (peerConnectionState) { @@ -299,6 +299,7 @@ internal class ClientEventReporter( joinReason, role, iceState, + wasPrevConnected, peerConnectionState, ) } @@ -347,9 +348,9 @@ internal class ClientEventReporter( joinReason: JoinReason, role: PeerConnectionRole, iceState: VideoAnalyticsIceState, + wasPrevConnected: Boolean, peerConnectionState: PeerConnection.PeerConnectionState, ) { - val wasPrevConnected = pcEverConnected[role] != null val stageId = UUID.randomUUID().toString() val now = System.currentTimeMillis() postCallFlightSessions[stageId] = PostCallFlightSession( @@ -400,7 +401,6 @@ internal class ClientEventReporter( peerConnectionState: PeerConnection.PeerConnectionState, ) { - pcEverConnected[role] = PcConnected(System.currentTimeMillis()) val pcState = pcEventReporterStateHolder.map.remove(peerConnectionHashCode) ?: return val stageId = pcState.stageId @@ -630,4 +630,3 @@ internal class PeerConnectionEventReporterState( var stageId: String, val peerConnectionRole: PeerConnectionRole, ) -internal class PcConnected(val lastConnectedTime: Long) diff --git a/stream-video-android-core/src/test/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalyticsTest.kt b/stream-video-android-core/src/test/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalyticsTest.kt index 4f74c1e8d25..6a45d9388c3 100644 --- a/stream-video-android-core/src/test/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalyticsTest.kt +++ b/stream-video-android-core/src/test/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalyticsTest.kt @@ -103,6 +103,7 @@ class PeerConnectionAnalyticsTest { joinReason = JoinReason.ReJoin, role = PeerConnectionRole.SUBSCRIBE, iceState = VideoAnalyticsIceState.NOT_CONNECTED, + wasPrevConnected = false, peerConnectionState = PeerConnection.PeerConnectionState.CONNECTING, ) } @@ -131,6 +132,7 @@ class PeerConnectionAnalyticsTest { joinReason = any(), role = PeerConnectionRole.PUBLISH, iceState = VideoAnalyticsIceState.CONNECTED, + wasPrevConnected = false, peerConnectionState = PeerConnection.PeerConnectionState.CONNECTING, ) } @@ -161,6 +163,7 @@ class PeerConnectionAnalyticsTest { joinReason = any(), role = PeerConnectionRole.PUBLISH, iceState = VideoAnalyticsIceState.NOT_CONNECTED, + wasPrevConnected = true, peerConnectionState = PeerConnection.PeerConnectionState.CONNECTED, ) } @@ -191,6 +194,7 @@ class PeerConnectionAnalyticsTest { joinReason = any(), role = PeerConnectionRole.PUBLISH, iceState = VideoAnalyticsIceState.NOT_CONNECTED, + wasPrevConnected = false, peerConnectionState = PeerConnection.PeerConnectionState.CONNECTING, ) } @@ -220,6 +224,7 @@ class PeerConnectionAnalyticsTest { joinReason = any(), role = PeerConnectionRole.PUBLISH, iceState = VideoAnalyticsIceState.FAILED, + wasPrevConnected = false, peerConnectionState = PeerConnection.PeerConnectionState.FAILED, ) } @@ -285,12 +290,72 @@ class PeerConnectionAnalyticsTest { joinReason = any(), role = PeerConnectionRole.PUBLISH, iceState = VideoAnalyticsIceState.CONNECTED, + wasPrevConnected = true, peerConnectionState = PeerConnection.PeerConnectionState.CONNECTED, ) } scope.cancel() } + @Test + fun `a publisher reconnect is flagged as previously connected without affecting subscriber`() { + val scope = CoroutineScope(Dispatchers.Unconfined) + val peerConnectionAnalytics = analytics(scope) + + peerConnectionAnalytics.onPeerConnectionStateChanged( + peerConnectionHashCode = 42, + role = PeerConnectionRole.PUBLISH, + iceState = VideoAnalyticsIceState.CONNECTED, + peerConnectionState = PeerConnection.PeerConnectionState.CONNECTED, + ) + + peerConnectionAnalytics.onPeerConnectionStateChanged( + peerConnectionHashCode = 43, + role = PeerConnectionRole.SUBSCRIBE, + iceState = VideoAnalyticsIceState.NOT_CONNECTED, + peerConnectionState = PeerConnection.PeerConnectionState.CONNECTING, + ) + + peerConnectionAnalytics.onPeerConnectionStateChanged( + peerConnectionHashCode = 44, + role = PeerConnectionRole.PUBLISH, + iceState = VideoAnalyticsIceState.NOT_CONNECTED, + peerConnectionState = PeerConnection.PeerConnectionState.CONNECTING, + ) + + verify { + reporter.onPeerConnectionStateChanged( + peerConnectionHashCode = 43, + callId = "call-1", + callType = "default", + joinStageAttemptId = any(), + callSessionId = any(), + sfuId = any(), + joinReason = any(), + role = PeerConnectionRole.SUBSCRIBE, + iceState = VideoAnalyticsIceState.NOT_CONNECTED, + wasPrevConnected = false, + peerConnectionState = PeerConnection.PeerConnectionState.CONNECTING, + ) + } + verify { + reporter.onPeerConnectionStateChanged( + peerConnectionHashCode = 44, + callId = "call-1", + callType = "default", + joinStageAttemptId = any(), + callSessionId = any(), + sfuId = any(), + joinReason = any(), + role = PeerConnectionRole.PUBLISH, + iceState = VideoAnalyticsIceState.NOT_CONNECTED, + wasPrevConnected = true, + peerConnectionState = PeerConnection.PeerConnectionState.CONNECTING, + ) + } + scope.cancel() + } + @Test fun `a connected publisher resets the stage back to completed`() = runTest { val session = mockSession( diff --git a/stream-video-android-core/src/test/kotlin/io/getstream/video/android/core/analytics/reporting/ClientEventReporterTest.kt b/stream-video-android-core/src/test/kotlin/io/getstream/video/android/core/analytics/reporting/ClientEventReporterTest.kt index 485ba640d24..1d7385fadaa 100644 --- a/stream-video-android-core/src/test/kotlin/io/getstream/video/android/core/analytics/reporting/ClientEventReporterTest.kt +++ b/stream-video-android-core/src/test/kotlin/io/getstream/video/android/core/analytics/reporting/ClientEventReporterTest.kt @@ -87,6 +87,7 @@ class ClientEventReporterTest { iceState: VideoAnalyticsIceState, role: PeerConnectionRole = PeerConnectionRole.PUBLISH, pcHashCode: Int = 100, + wasPrevConnected: Boolean = false, ) = reporter.onPeerConnectionStateChanged( peerConnectionHashCode = pcHashCode, callId = "call-1", @@ -97,6 +98,7 @@ class ClientEventReporterTest { joinReason = JoinReason.FirstAttempt, role = role, iceState = iceState, + wasPrevConnected = wasPrevConnected, peerConnectionState = pcState, ) @@ -329,6 +331,7 @@ class ClientEventReporterTest { pcState = PeerConnection.PeerConnectionState.CONNECTING, iceState = VideoAnalyticsIceState.NOT_CONNECTED, pcHashCode = 2, + wasPrevConnected = true, ) val reconnectInitiated = dispatcher.sent.last() From ec92c1c778456eb0b98c5b0f90e744c0843f5fea Mon Sep 17 00:00:00 2001 From: rahullohra Date: Fri, 14 Aug 2026 17:15:03 +0530 Subject: [PATCH 2/2] test: add remaining unit-tests --- .../call/observer/PeerConnectionAnalytics.kt | 20 +++---- .../PeerConnectionAnalyticsStateHolder.kt | 8 +++ .../observer/PeerConnectionAnalyticsTest.kt | 57 ++++++++++--------- 3 files changed, 48 insertions(+), 37 deletions(-) diff --git a/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalytics.kt b/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalytics.kt index e6f064cc290..0d7fbf36f02 100644 --- a/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalytics.kt +++ b/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalytics.kt @@ -162,16 +162,7 @@ internal class PeerConnectionAnalytics( iceState: VideoAnalyticsIceState, peerConnectionState: PeerConnection.PeerConnectionState?, ) { - when (peerConnectionState) { - PeerConnection.PeerConnectionState.CONNECTED -> { - if (role == PeerConnectionRole.PUBLISH) { - stateHolder.updatePublisherEverConnected(true) - } else { - stateHolder.updateSubscriberEverConnected(true) - } - } - else -> {} - } + val wasPrevConnected = stateHolder.isPcEverConnected(role) reporter.onPeerConnectionStateChanged( peerConnectionHashCode = peerConnectionHashCode, @@ -183,9 +174,16 @@ internal class PeerConnectionAnalytics( joinStageAttemptId = joinAnalyticsStateHolder.state.value.joinStageAttemptId ?: "unknown", joinReason = joinAnalyticsStateHolder.state.value.joinReason ?: JoinReason.Unknown, sfuId = sfuAnalyticsStateHolder.sfuId.value, - wasPrevConnected = stateHolder.isPcEverConnected(role), + wasPrevConnected = wasPrevConnected, callSessionId = joinAnalyticsStateHolder.state.value.callSessionId, ) + + when (peerConnectionState) { + PeerConnection.PeerConnectionState.CONNECTED -> { + stateHolder.setPcEverConnected(role, true) + } + else -> {} + } } fun stop() { diff --git a/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalyticsStateHolder.kt b/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalyticsStateHolder.kt index b51e392007a..21c665d635c 100644 --- a/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalyticsStateHolder.kt +++ b/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalyticsStateHolder.kt @@ -64,6 +64,14 @@ internal class PeerConnectionAnalyticsStateHolder { } } + fun setPcEverConnected(role: PeerConnectionRole, connected: Boolean) { + if (role == PeerConnectionRole.PUBLISH) { + updatePublisherEverConnected(connected) + } else { + updateSubscriberEverConnected(connected) + } + } + fun update( peerConnectionObserverJob: Job? = state.value.peerConnectionObserverJob, publisherJob: Job? = state.value.publisherJob, diff --git a/stream-video-android-core/src/test/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalyticsTest.kt b/stream-video-android-core/src/test/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalyticsTest.kt index 6a45d9388c3..e6fda59a1cd 100644 --- a/stream-video-android-core/src/test/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalyticsTest.kt +++ b/stream-video-android-core/src/test/kotlin/io/getstream/video/android/core/analytics/call/observer/PeerConnectionAnalyticsTest.kt @@ -109,6 +109,32 @@ class PeerConnectionAnalyticsTest { } } + @Test + fun `a new call starts with wasPrevConnected set to false`() { + analytics(CoroutineScope(Dispatchers.Unconfined)).onPeerConnectionStateChanged( + peerConnectionHashCode = 42, + role = PeerConnectionRole.PUBLISH, + iceState = VideoAnalyticsIceState.NOT_CONNECTED, + peerConnectionState = PeerConnection.PeerConnectionState.CONNECTING, + ) + + verify(exactly = 1) { + reporter.onPeerConnectionStateChanged( + peerConnectionHashCode = 42, + callId = "call-1", + callType = "default", + joinStageAttemptId = any(), + callSessionId = any(), + sfuId = any(), + joinReason = any(), + role = PeerConnectionRole.PUBLISH, + iceState = VideoAnalyticsIceState.NOT_CONNECTED, + wasPrevConnected = false, + peerConnectionState = PeerConnection.PeerConnectionState.CONNECTING, + ) + } + } + @Test fun `a connecting publisher reports its current ice state immediately and marks the stage in progress`() = runTest { val session = mockSession( @@ -163,7 +189,7 @@ class PeerConnectionAnalyticsTest { joinReason = any(), role = PeerConnectionRole.PUBLISH, iceState = VideoAnalyticsIceState.NOT_CONNECTED, - wasPrevConnected = true, + wasPrevConnected = false, peerConnectionState = PeerConnection.PeerConnectionState.CONNECTED, ) } @@ -290,15 +316,16 @@ class PeerConnectionAnalyticsTest { joinReason = any(), role = PeerConnectionRole.PUBLISH, iceState = VideoAnalyticsIceState.CONNECTED, - wasPrevConnected = true, + wasPrevConnected = false, peerConnectionState = PeerConnection.PeerConnectionState.CONNECTED, ) } + assertTrue(stateHolder.isPcEverConnected(PeerConnectionRole.PUBLISH)) scope.cancel() } @Test - fun `a publisher reconnect is flagged as previously connected without affecting subscriber`() { + fun `an existing call whose Publisher was previously connected sends wasPrevConnected as true`() { val scope = CoroutineScope(Dispatchers.Unconfined) val peerConnectionAnalytics = analytics(scope) @@ -311,19 +338,12 @@ class PeerConnectionAnalyticsTest { peerConnectionAnalytics.onPeerConnectionStateChanged( peerConnectionHashCode = 43, - role = PeerConnectionRole.SUBSCRIBE, - iceState = VideoAnalyticsIceState.NOT_CONNECTED, - peerConnectionState = PeerConnection.PeerConnectionState.CONNECTING, - ) - - peerConnectionAnalytics.onPeerConnectionStateChanged( - peerConnectionHashCode = 44, role = PeerConnectionRole.PUBLISH, iceState = VideoAnalyticsIceState.NOT_CONNECTED, peerConnectionState = PeerConnection.PeerConnectionState.CONNECTING, ) - verify { + verify(exactly = 1) { reporter.onPeerConnectionStateChanged( peerConnectionHashCode = 43, callId = "call-1", @@ -332,21 +352,6 @@ class PeerConnectionAnalyticsTest { callSessionId = any(), sfuId = any(), joinReason = any(), - role = PeerConnectionRole.SUBSCRIBE, - iceState = VideoAnalyticsIceState.NOT_CONNECTED, - wasPrevConnected = false, - peerConnectionState = PeerConnection.PeerConnectionState.CONNECTING, - ) - } - verify { - reporter.onPeerConnectionStateChanged( - peerConnectionHashCode = 44, - callId = "call-1", - callType = "default", - joinStageAttemptId = any(), - callSessionId = any(), - sfuId = any(), - joinReason = any(), role = PeerConnectionRole.PUBLISH, iceState = VideoAnalyticsIceState.NOT_CONNECTED, wasPrevConnected = true,