From ae06ad83726dc272f1a61626813838f4d81b26ce Mon Sep 17 00:00:00 2001 From: pratimmallick Date: Tue, 11 Aug 2026 16:00:32 +0530 Subject: [PATCH 1/2] fix(core): always announce video layers on SetPublisher Avoid empty-layer announces when the track is not LIVE yet, which caused ERROR_CODE_REQUEST_VALIDATION_FAILED and an unnecessary full rejoin. Co-authored-by: Cursor --- .../android/core/call/connection/Publisher.kt | 45 +++++++----- .../core/call/connection/PublisherTest.kt | 71 ++++++++++++++++++- 2 files changed, 97 insertions(+), 19 deletions(-) diff --git a/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/call/connection/Publisher.kt b/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/call/connection/Publisher.kt index 29b9cfbbe8..e94e508a97 100644 --- a/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/call/connection/Publisher.kt +++ b/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/call/connection/Publisher.kt @@ -206,25 +206,36 @@ internal class Publisher( session_id = sessionId, ) val response = sfuClient.setPublisher(request) - logger.i { "Received answer: ${response.sdp}" } if (response.error != null) { logger.e { "SetPublisherRequest Received error: ${response.error}, SetPublisherRequest: $request" } tracer.trace("negotiate-error-setpublisher", response.error.message ?: "unknown") - logger.e { "rejoin cause error in sfuClient.setPublisher, message:${response.error.message}" } - when (response.error.code) { /** - * We are getting this error right away after joining the call first time - * Full error: 16:04:05.032 Call:PeerC...:publisher E (DefaultDispatcher-worker-17:526) SetPublisherRequest Received error: Error{code=ERROR_CODE_REQUEST_VALIDATION_FAILED, message=Invalid SetPublisher request, should_retry=false} - * This will cause emission of ParticipantLeftEvent to other person + * Historically common right after first join when video layers were omitted + * for non-LIVE tracks. Layers are now always computed from publish options, so + * remaining validation failures are treated as a real publisher/SFU mismatch. */ - ErrorCode.ERROR_CODE_REQUEST_VALIDATION_FAILED -> rejoin() + ErrorCode.ERROR_CODE_REQUEST_VALIDATION_FAILED -> { + logger.e { + "rejoin cause error in sfuClient.setPublisher, " + + "message:${response.error.message}" + } + rejoin() + } - else -> {} + else -> { + logger.e { + "Unhandled SetPublisher error code=${response.error.code}, " + + "message=${response.error.message}" + } + } } + return@submit } + + logger.i { "Received answer: ${response.sdp}" } setRemoteDescription(SessionDescription(SessionDescription.Type.ANSWER, response.sdp)) .onErrorSuspend { tracer.trace( @@ -234,7 +245,6 @@ internal class Publisher( }.onSuccess { logger.d { "Publisher negotiation successfully done ✅" } } - // Set ice trickle } isIceRestarting = false } @@ -649,16 +659,15 @@ internal class Publisher( } val isTrackLive = track.state() == MediaStreamTrack.State.LIVE val isAudio = isAudioTrackType(publishOption.track_type) + // Layer math only needs dimension + PublishOption — not LIVE/frames. Previously we skipped + // computeLayers when !LIVE, announced empty layers, and SFU rejected SetPublisher on first + // join (then we rejoined). Always compute so muted/non-LIVE video still has layers. val layers = if (!isAudio) { - if (isTrackLive) { - computeLayers( - captureFormat, - track, - publishOption, - ) - } else { - transceiverCache.getLayers(publishOption) - } + computeLayers( + captureFormat, + track, + publishOption, + ) ?: transceiverCache.getLayers(publishOption) } else { null } diff --git a/stream-video-android-core/src/test/kotlin/io/getstream/video/android/core/call/connection/PublisherTest.kt b/stream-video-android-core/src/test/kotlin/io/getstream/video/android/core/call/connection/PublisherTest.kt index 326221df03..cd3245c4b3 100644 --- a/stream-video-android-core/src/test/kotlin/io/getstream/video/android/core/call/connection/PublisherTest.kt +++ b/stream-video-android-core/src/test/kotlin/io/getstream/video/android/core/call/connection/PublisherTest.kt @@ -62,7 +62,10 @@ import stream.video.sfu.event.VideoSender import stream.video.sfu.models.AudioBitrateProfile import stream.video.sfu.models.Codec import stream.video.sfu.models.DegradationPreference +import stream.video.sfu.models.Error +import stream.video.sfu.models.ErrorCode import stream.video.sfu.models.PublishOption +import stream.video.sfu.models.TrackInfo import stream.video.sfu.models.TrackType import stream.video.sfu.models.VideoDimension import stream.video.sfu.signal.SetPublisherResponse @@ -92,6 +95,7 @@ class PublisherTest { private lateinit var publisher: Publisher private val coroutineContext = UnconfinedTestDispatcher() private val testScope = TestScope(coroutineContext) + private var rejoinInvocations = 0 //region Example PublishOptions private val videoPublishOption = PublishOption( @@ -129,6 +133,7 @@ class PublisherTest { @Before fun setUp() { MockKAnnotations.init(this, relaxUnitFun = true) + rejoinInvocations = 0 // Mock the mediaManager and peerConnectionFactory so they return mock Audio/Video tracks. every { mockPeerConnectionFactory.makeAudioTrack(any(), any()) } answers { @@ -171,7 +176,7 @@ class PublisherTest { maxBitRate = 1_500_000, sfuClient = mockSignalServerService, sessionId = "session-id", - rejoin = { }, + rejoin = { rejoinInvocations++ }, tracer = mockk(relaxed = true), fastReconnect = {}, transceiverCache = mockTransceiverCache, @@ -285,6 +290,70 @@ class PublisherTest { coVerify(exactly = 0) { mockSignalServerService.setPublisher(any()) } } + @Test + fun `getAnnouncedTracks includes video layers even when track is not LIVE`() { + val mockVideoTrack = mockk(relaxed = true) { + every { id() } returns "video-1" + every { kind() } returns "video" + every { isDisposed } returns false + every { state() } returns MediaStreamTrack.State.ENDED + } + val mockSender = mockk(relaxed = true) { + every { track() } returns mockVideoTrack + } + val mockTransceiver = mockk(relaxed = true) { + every { sender } returns mockSender + every { mid } returns "0" + } + every { mockTransceiverCache.items() } returns listOf( + TransceiverId(videoPublishOption, mockTransceiver), + ) + every { mockTransceiverCache.indexOf(videoPublishOption) } returns 0 + every { mockTransceiverCache.getLayers(videoPublishOption) } returns null + + val announced = publisher.getAnnouncedTracks(null, null) + + assertEquals(1, announced.size) + assertTrue(announced[0].muted) + assertTrue(announced[0].layers.isNotEmpty()) + } + + @Test + fun `SetPublisher validation failure rejoins`() = runTest(coroutineContext) { + every { publisher.getAnnouncedTracks(any(), any()) } returns listOf( + TrackInfo( + track_id = "video-1", + track_type = TrackType.TRACK_TYPE_VIDEO, + mid = "0", + muted = false, + layers = listOf( + stream.video.sfu.models.VideoLayer( + rid = "f", + video_dimension = VideoDimension(1280, 720), + bitrate = 1_000_000, + fps = 30, + ), + ), + publish_option_id = videoPublishOption.id, + ), + ) + coEvery { publisher.setLocalDescription(any()) } returns Result.Success(Unit) + coEvery { publisher.setRemoteDescription(any()) } returns Result.Success(Unit) + coEvery { mockSignalServerService.setPublisher(any()) } returns SetPublisherResponse( + sdp = "", + error = Error( + code = ErrorCode.ERROR_CODE_REQUEST_VALIDATION_FAILED, + message = "Invalid SetPublisher request", + should_retry = false, + ), + ) + + publisher.negotiate(source = "test") + + coVerify(exactly = 1) { mockSignalServerService.setPublisher(any()) } + assertEquals(1, rejoinInvocations) + } + @Test fun `close with stopTracks = true stops publishing and closes connection`() = runTest { val mockVideoTrack = mockk(relaxed = true) { From 016d192fb776b5846bf6a2c5cd379f159f779b35 Mon Sep 17 00:00:00 2001 From: pratimmallick Date: Mon, 17 Aug 2026 16:23:13 +0530 Subject: [PATCH 2/2] fix(core): clear isIceRestarting on SetPublisher error paths Early return@submit inside negotiate skipped the flag reset and could stall later renegotiation after an iceRestart failure. Co-authored-by: Cursor --- .../android/core/call/connection/Publisher.kt | 10 ++-- .../core/call/connection/PublisherTest.kt | 48 +++++++++++++++++++ 2 files changed, 55 insertions(+), 3 deletions(-) diff --git a/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/call/connection/Publisher.kt b/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/call/connection/Publisher.kt index e94e508a97..09134db257 100644 --- a/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/call/connection/Publisher.kt +++ b/stream-video-android-core/src/main/kotlin/io/getstream/video/android/core/call/connection/Publisher.kt @@ -195,8 +195,8 @@ internal class Publisher( logger.i { "Negotiating with tracks: $trackInfos" } logger.i { "Offer: ${offer.description}" } - safeCall { - isIceRestarting = iceRestart + isIceRestarting = iceRestart + try { setLocalDescription(offer).onErrorSuspend { tracer.trace("negotiate-error-setlocaldescription", it.message ?: "unknown") } @@ -245,8 +245,12 @@ internal class Publisher( }.onSuccess { logger.d { "Publisher negotiation successfully done ✅" } } + } catch (e: Exception) { + logger.e(e) { "[negotiate] Exception occurred: ${e.message}" } + } finally { + // Must clear even when returning early from @submit (inline safeCall used to skip this). + isIceRestarting = false } - isIceRestarting = false } override suspend fun stats(): ComputedStats? = safeCallWithDefault(null) { diff --git a/stream-video-android-core/src/test/kotlin/io/getstream/video/android/core/call/connection/PublisherTest.kt b/stream-video-android-core/src/test/kotlin/io/getstream/video/android/core/call/connection/PublisherTest.kt index cd3245c4b3..04ea80f295 100644 --- a/stream-video-android-core/src/test/kotlin/io/getstream/video/android/core/call/connection/PublisherTest.kt +++ b/stream-video-android-core/src/test/kotlin/io/getstream/video/android/core/call/connection/PublisherTest.kt @@ -351,9 +351,57 @@ class PublisherTest { publisher.negotiate(source = "test") coVerify(exactly = 1) { mockSignalServerService.setPublisher(any()) } + coVerify(exactly = 0) { publisher.setRemoteDescription(any()) } assertEquals(1, rejoinInvocations) } + @Test + fun `iceRestart SetPublisher error clears isIceRestarting so later negotiate runs`() = runTest( + coroutineContext, + ) { + every { publisher.getAnnouncedTracks(any(), any()) } returns listOf( + TrackInfo( + track_id = "video-1", + track_type = TrackType.TRACK_TYPE_VIDEO, + mid = "0", + muted = false, + layers = listOf( + stream.video.sfu.models.VideoLayer( + rid = "f", + video_dimension = VideoDimension(1280, 720), + bitrate = 1_000_000, + fps = 30, + ), + ), + publish_option_id = videoPublishOption.id, + ), + ) + coEvery { publisher.setLocalDescription(any()) } returns Result.Success(Unit) + coEvery { publisher.setRemoteDescription(any()) } returns Result.Success(Unit) + coEvery { mockSignalServerService.setPublisher(any()) } returnsMany listOf( + SetPublisherResponse( + sdp = "", + error = Error( + code = ErrorCode.ERROR_CODE_PARTICIPANT_NOT_FOUND, + message = "participant not found", + should_retry = false, + ), + ), + SetPublisherResponse( + sdp = fakeSdpAnswer.description, + error = null, + ), + ) + + publisher.negotiate(source = "ice-restart", iceRestart = true) + assertEquals(0, rejoinInvocations) + + // Without finally clearing isIceRestarting, this second call would no-op. + publisher.negotiate(source = "after-ice-restart-error") + + coVerify(exactly = 2) { mockSignalServerService.setPublisher(any()) } + } + @Test fun `close with stopTracks = true stops publishing and closes connection`() = runTest { val mockVideoTrack = mockk(relaxed = true) {