Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/fix-reconnect-publication-job-leak.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"client-sdk-android": patch
---

Fixed local track publications leaking their jobs on every full reconnect, which left the audio feature collectors of the old publications running and sending feature updates for stale track sids.
12 changes: 6 additions & 6 deletions livekit-android-sdk/detekt-baseline-release.xml
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@
<ID>ComplexCondition:RemoteTrackPublication.kt$RemoteTrackPublication$isAutoManaged || !subscribed || this.fps == fps || track !is VideoTrack</ID>
<ID>ComplexCondition:RemoteTrackPublication.kt$RemoteTrackPublication$isAutoManaged || !subscribed || videoDimensions == dimensions || track !is VideoTrack</ID>
<ID>ComplexCondition:TextureViewRenderer.kt$TextureViewRenderer$enableFixedSize &amp;&amp; rotatedFrameWidth != 0 &amp;&amp; rotatedFrameHeight != 0 &amp;&amp; width != 0 &amp;&amp; height != 0</ID>
<ID>CyclomaticComplexMethod:LocalParticipant.kt$LocalParticipant$@Throws(TrackException.PublishException::class) private suspend fun publishTrackImpl( track: Track, options: TrackPublishOptions, requestConfig: AddTrackRequest.Builder.() -> Unit, encodings: List&lt;RtpParameters.Encoding> = emptyList(), publishListener: PublishListener? = null, ): LocalTrackPublication?</ID>
<ID>CyclomaticComplexMethod:LocalParticipant.kt$LocalParticipant$@Throws(TrackException.PublishException::class) private suspend fun publishTrackImpl( track: Track, options: TrackPublishOptions, requestConfig: AddTrackRequest.Builder.() -> Unit, encodings: List&lt;RtpParameters.Encoding> = emptyList(), publishListener: PublishListener? = null, ): PublishResult?</ID>
<ID>CyclomaticComplexMethod:LocalParticipant.kt$LocalParticipant$private fun computeVideoEncodings( isScreenShare: Boolean, dimensions: Track.Dimensions, options: VideoTrackPublishOptions, ): List&lt;RtpParameters.Encoding></ID>
<ID>CyclomaticComplexMethod:LocalParticipant.kt$LocalParticipant$private suspend fun setTrackEnabled( source: Track.Source, enabled: Boolean, screenCaptureParams: ScreenCaptureParams? = null, ): Boolean</ID>
<ID>CyclomaticComplexMethod:LocalParticipant.kt$LocalParticipant$suspend fun publishVideoTrack( track: LocalVideoTrack, options: VideoTrackPublishOptions = VideoTrackPublishOptions( null, if (track.options.isScreencast) screenShareTrackPublishDefaults else videoTrackPublishDefaults, ), publishListener: PublishListener? = null, ): Boolean</ID>
Expand All @@ -27,13 +27,13 @@
<ID>CyclomaticComplexMethod:Room.kt$Room$@Throws(Exception::class) suspend fun connect(url: String, token: String, options: ConnectOptions = ConnectOptions())</ID>
<ID>CyclomaticComplexMethod:RoomEvent.kt$fun LivekitModels.DisconnectReason?.convert(): DisconnectReason</ID>
<ID>CyclomaticComplexMethod:SignalClient.kt$SignalClient$override fun onFailure(webSocket: WebSocket, t: Throwable, response: Response?)</ID>
<ID>CyclomaticComplexMethod:SignalClient.kt$SignalClient$private fun handleSignalResponseImpl(ws: WebSocket, response: LivekitRtc.SignalResponse)</ID>
<ID>CyclomaticComplexMethod:SignalClient.kt$SignalClient$private fun handleSignalResponseImpl(connection: SignalConnection, response: LivekitRtc.SignalResponse)</ID>
<ID>EmptyFunctionBlock:RTCEngine.kt$RTCEngine${ }</ID>
<ID>HasPlatformType:DataChannelManager.kt$DataChannelManager$@get:FlowObservable var state by flowDelegate(dataChannel.state()) private set</ID>
<ID>IgnoredReturnValue:RpcServerManager.kt$RpcServerManager$publishRpcAck(callerIdentity, requestId)</ID>
<ID>InstanceOfCheckForException:RpcServerManager.kt$RpcServerManager$e is RpcError</ID>
<ID>LargeClass:LocalParticipant.kt$LocalParticipant : ParticipantOutgoingDataStreamManagerRpcManager</ID>
<ID>LargeClass:RTCEngine.kt$RTCEngine : Listener</ID>
<ID>LargeClass:RTCEngine.kt$RTCEngine : ListenerSignalSessionListener</ID>
<ID>LargeClass:Room.kt$Room : ListenerParticipantListenerRpcManagerIncomingDataStreamManager</ID>
<ID>LargeClass:SignalClient.kt$SignalClient : WebSocketListener</ID>
<ID>LongMethod:RTCEngine.kt$RTCEngine$@Synchronized @VisibleForTesting(otherwise = VisibleForTesting.PACKAGE_PRIVATE) fun reconnect()</ID>
Expand Down Expand Up @@ -61,7 +61,7 @@
<ID>LongParameterList:Room.kt$Room$( @Assisted private val context: Context, internal val engine: RTCEngine, private val eglBase: EglBase, localParticipantFactory: LocalParticipant.Factory, private val defaultsManager: DefaultsManager, @Named(InjectionNames.DISPATCHER_DEFAULT) private val defaultDispatcher: CoroutineDispatcher, @Named(InjectionNames.DISPATCHER_IO) private val ioDispatcher: CoroutineDispatcher, /** * The [AudioHandler] for setting up the audio as need. * * By default, this is an instance of [AudioSwitchHandler]. * * This can be substituted for your own custom implementation through * [LiveKitOverrides.audioOptions] when creating the room with [LiveKit.create]. * * @see [audioSwitchHandler] * @see [AudioSwitchHandler] */ val audioHandler: AudioHandler, private val closeableManager: CloseableManager, private val e2EEManagerFactory: E2EEManager.Factory, private val communicationWorkaround: CommunicationWorkaround, val audioProcessingController: AudioProcessingController, /** * A holder for objects that are used internally within LiveKit. */ val lkObjects: LKObjects, networkCallbackManagerFactory: NetworkCallbackManagerFactory, private val audioDeviceModule: AudioDeviceModule, private val regionUrlProviderFactory: RegionUrlProvider.Factory, private val connectionWarmer: ConnectionWarmer, private val audioRecordPrewarmer: AudioRecordPrewarmer, private val incomingDataStreamManager: IncomingDataStreamManager, private val rpcClientManager: RpcClientManager, private val rpcServerManager: RpcServerManager, private val remoteParticipantFactory: RemoteParticipant.Factory, )</ID>
<ID>MapGetWithNotNullAssertionOperator:LocalParticipant.kt$LocalParticipant$sourcePubLocks[source]!!</ID>
<ID>NestedBlockDepth:ByteStreamSender.kt$@CheckResult suspend fun ByteStreamSender.write(source: Source): Result&lt;Unit></ID>
<ID>NestedBlockDepth:LocalParticipant.kt$LocalParticipant$@Throws(TrackException.PublishException::class) private suspend fun publishTrackImpl( track: Track, options: TrackPublishOptions, requestConfig: AddTrackRequest.Builder.() -> Unit, encodings: List&lt;RtpParameters.Encoding> = emptyList(), publishListener: PublishListener? = null, ): LocalTrackPublication?</ID>
<ID>NestedBlockDepth:LocalParticipant.kt$LocalParticipant$@Throws(TrackException.PublishException::class) private suspend fun publishTrackImpl( track: Track, options: TrackPublishOptions, requestConfig: AddTrackRequest.Builder.() -> Unit, encodings: List&lt;RtpParameters.Encoding> = emptyList(), publishListener: PublishListener? = null, ): PublishResult?</ID>
<ID>NestedBlockDepth:LocalParticipant.kt$LocalParticipant$fun cleanup()</ID>
<ID>NestedBlockDepth:LocalVideoTrack.kt$LocalVideoTrack$internal fun setPublishingCodecs(codecs: List&lt;SubscribedCodec>): List&lt;VideoCodec></ID>
<ID>NestedBlockDepth:LocalVideoTrack.kt$LocalVideoTrack$private fun setPublishingLayersForSender( sender: RtpSender, qualities: List&lt;LivekitRtc.SubscribedQuality>, )</ID>
Expand All @@ -71,7 +71,7 @@
<ID>NestedBlockDepth:RTCEngine.kt$RTCEngine$private fun makeRTCConfig( serverResponse: Either&lt;JoinResponse, ReconnectResponse>, connectOptions: ConnectOptions, ): RTCConfiguration</ID>
<ID>NestedBlockDepth:Room.kt$Room$override suspend fun onPostReconnect(isFullReconnect: Boolean)</ID>
<ID>NestedBlockDepth:SignalClient.kt$SignalClient$override fun onFailure(webSocket: WebSocket, t: Throwable, response: Response?)</ID>
<ID>NestedBlockDepth:SignalClient.kt$SignalClient$private fun handleSignalResponse(ws: WebSocket, response: LivekitRtc.SignalResponse)</ID>
<ID>NestedBlockDepth:SignalClient.kt$SignalClient$private fun handleSignalResponse(connection: SignalConnection, response: LivekitRtc.SignalResponse)</ID>
<ID>SwallowedException:FlowExt.kt$e: CancellationException</ID>
<ID>SwallowedException:LocalVideoTrack.kt$LocalVideoTrack$e: Exception</ID>
<ID>SwallowedException:TextureViewRenderer.kt$TextureViewRenderer$e: NotFoundException</ID>
Expand All @@ -86,7 +86,7 @@
<ID>TooManyFunctions:Participant.kt$ParticipantListener</ID>
<ID>TooManyFunctions:PeerConnectionTransport.kt$PeerConnectionTransport</ID>
<ID>TooManyFunctions:PublisherTransportObserver.kt$PublisherTransportObserver : ObserverListenerPeerConnectionStateObservable</ID>
<ID>TooManyFunctions:RTCEngine.kt$RTCEngine : Listener</ID>
<ID>TooManyFunctions:RTCEngine.kt$RTCEngine : ListenerSignalSessionListener</ID>
<ID>TooManyFunctions:RTCEngine.kt$RTCEngine$Listener</ID>
<ID>TooManyFunctions:RTCMetricsManager.kt$io.livekit.android.room.metrics.RTCMetricsManager.kt</ID>
<ID>TooManyFunctions:RTCModule.kt$RTCModule</ID>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,7 @@ import livekit.org.webrtc.RtpTransceiver
import livekit.org.webrtc.RtpTransceiver.RtpTransceiverInit
import livekit.org.webrtc.SessionDescription
import java.nio.ByteBuffer
import java.util.concurrent.atomic.AtomicLong
import javax.inject.Inject
import javax.inject.Named
import javax.inject.Singleton
Expand All @@ -119,7 +120,8 @@ internal constructor(
private val ioDispatcher: CoroutineDispatcher,
private val rtcThreadToken: RTCThreadToken,
private val dataPacketCryptorFactory: DataPacketCryptorManager.Factory,
) : SignalClient.Listener {
) : SignalClient.Listener,
SignalClient.SignalSessionListener {
internal var listener: Listener? = null

/**
Expand Down Expand Up @@ -175,9 +177,20 @@ internal constructor(

internal var reconnectPolicy: ReconnectPolicy = DefaultReconnectPolicy()

private val pendingTrackResolvers: MutableMap<String, Continuation<LivekitModels.TrackInfo>> =
private val pendingTrackResolvers: MutableMap<String, Continuation<PublishAcceptance>> =
mutableMapOf()

// Counts full reconnect preparations, which invalidate the server-side state of every
// publish accepted before them. Soft reconnects preserve publishes and do not advance
// it.
private val fullReconnectEpoch = AtomicLong(0)

internal fun currentFullReconnectEpoch(): Long = fullReconnectEpoch.get()

internal fun advanceFullReconnectEpoch() {
fullReconnectEpoch.incrementAndGet()
}

internal var regionUrlProvider: RegionUrlProvider? = null
private var sessionUrl: String? = null
private var sessionToken: String? = null
Expand Down Expand Up @@ -266,6 +279,7 @@ internal constructor(
if (connectionState == ConnectionState.DISCONNECTED) {
connectionState = ConnectionState.CONNECTING
}
client.prepareSignalConnection(fullReconnectEpoch.get())
val joinResponse = client.join(url, token, options, roomOptions)
ensureActive()

Expand Down Expand Up @@ -394,13 +408,13 @@ internal constructor(
/**
* @param builder an optional builder to include other parameters related to the track
*/
suspend fun addTrack(
internal suspend fun addTrack(
cid: String,
name: String,
kind: LivekitModels.TrackType,
stream: String?,
builder: LivekitRtc.AddTrackRequest.Builder = LivekitRtc.AddTrackRequest.newBuilder(),
): LivekitModels.TrackInfo {
): PublishAcceptance {
synchronized(pendingTrackResolvers) {
if (pendingTrackResolvers[cid] != null) {
throw TrackException.DuplicateTrackException("Track with same ID $cid has already been published!")
Expand Down Expand Up @@ -639,6 +653,7 @@ internal constructor(
LKLog.v { "Attempting soft reconnect." }
subscriber?.prepareForIceRestart()
try {
client.prepareSignalConnection(fullReconnectEpoch.get())
val response = client.reconnect(url!!, token, participantSid)
if (response is Either.Left) {
val reconnectResponse = response.value
Expand Down Expand Up @@ -1193,6 +1208,14 @@ internal constructor(
}

override fun onLocalTrackPublished(response: LivekitRtc.TrackPublishedResponse) {
handleLocalTrackPublished(response, fullReconnectEpoch.get())
}

override fun onLocalTrackPublishedInSession(response: LivekitRtc.TrackPublishedResponse, fullReconnectEpoch: Long) {
handleLocalTrackPublished(response, fullReconnectEpoch)
}

private fun handleLocalTrackPublished(response: LivekitRtc.TrackPublishedResponse, fullReconnectEpoch: Long) {
val cid = response.cid ?: run {
LKLog.e { "local track published with null cid?" }
return
Expand All @@ -1211,7 +1234,7 @@ internal constructor(
LKLog.d { "missing track resolver for: $cid" }
return
}
cont.resume(response.track)
cont.resume(PublishAcceptance(response.track, fullReconnectEpoch))
}

override fun onLocalTrackSubscribed(trackSubscribed: LivekitRtc.TrackSubscribed) {
Expand Down Expand Up @@ -1607,6 +1630,15 @@ internal class SenderTransceiverHandle(
internal val signalSessionState: SignalSessionState,
)

/**
* A server-accepted publish: the track info from the TrackPublished response, stamped with
* the full reconnect epoch of the signal session that delivered it.
*/
internal class PublishAcceptance(
val trackInfo: LivekitModels.TrackInfo,
val fullReconnectEpoch: Long,
)

/**
* @suppress
*/
Expand Down
Loading
Loading