diff --git a/.changeset/client-telemetry.md b/.changeset/client-telemetry.md new file mode 100644 index 00000000..59f665ea --- /dev/null +++ b/.changeset/client-telemetry.md @@ -0,0 +1,5 @@ +--- +"client-sdk-android": minor +--- + +Client telemetry through the shared Rust core: each Room reports its connect, reconnect, publish and subscribe spans, RTC statistics, SDK warnings and errors and device state to its LiveKit Cloud project when the token carries the observability grant; apps can add `Room.emitTelemetryEvent(name, attributes)` and `Room.setTelemetryAttribute(key, value)`, and opt out with `LiveKit.disableTelemetry()`. diff --git a/.github/workflows/android.yml b/.github/workflows/android.yml index 39aa20eb..b601bda6 100644 --- a/.github/workflows/android.yml +++ b/.github/workflows/android.yml @@ -58,6 +58,28 @@ jobs: - name: Build and test with Gradle run: ./gradlew assembleRelease livekit-android-test:testRelease + # The telemetry end-to-end test runs the real Rust core under Robolectric: a host build of + # livekit-uniffi at the release gradle/libs.versions.toml pins, posting OTLP to a collector + # that writes what the test reads back (livekit-android-test/src/test/resources/telemetry/otelcol.yaml). + # Only this step sets LK_TELEMETRY_ENDPOINT, and a missing collector or library fails the job + # instead of skipping the test. The telemetry package's platform tests that need the core run here too. + - name: Telemetry E2E test + run: | + version="$(sed -n 's/^livekit-uniffi = "\(.*\)"/\1/p' gradle/libs.versions.toml)" + git clone --depth 1 --branch "livekit-uniffi/v$version" https://github.com/livekit/rust-sdks.git "$RUNNER_TEMP/rust-sdks" + (cd "$RUNNER_TEMP/rust-sdks" && cargo build -p livekit-uniffi) + curl -sSfL -o "$RUNNER_TEMP/otelcol.tar.gz" "https://github.com/open-telemetry/opentelemetry-collector-releases/releases/download/v0.162.0/otelcol-contrib_0.162.0_linux_amd64.tar.gz" + echo "fcc063749f730f8c21fe29f2d340ff174f5f1c5885bd3156fb6c985a3036fcc3 $RUNNER_TEMP/otelcol.tar.gz" | sha256sum -c - + tar -xzf "$RUNNER_TEMP/otelcol.tar.gz" -C "$RUNNER_TEMP" otelcol-contrib + nohup "$RUNNER_TEMP/otelcol-contrib" --config livekit-android-test/src/test/resources/telemetry/otelcol.yaml > "$RUNNER_TEMP/otelcol.log" 2>&1 & + timeout 30 bash -c 'until curl -s -o /dev/null http://127.0.0.1:4319; do sleep 1; done' || { cat "$RUNNER_TEMP/otelcol.log"; exit 1; } + ./gradlew livekit-android-test:testReleaseUnitTest --tests 'io.livekit.android.telemetry.*' --rerun \ + -PlivekitUniffiLibraryPath="$RUNNER_TEMP/rust-sdks/target/debug" + grep -q 'skipped="0"' livekit-android-test/build/test-results/testReleaseUnitTest/TEST-io.livekit.android.telemetry.TelemetryMockE2ETest.xml \ + || { echo "TelemetryMockE2ETest was skipped"; exit 1; } + env: + LK_TELEMETRY_ENDPOINT: http://127.0.0.1:4319 + - name: Run Detekt run: ./gradlew livekit-android-sdk:detektRelease diff --git a/gradle/libs.versions.toml b/gradle/libs.versions.toml index e0e6c0fd..50822f6f 100644 --- a/gradle/libs.versions.toml +++ b/gradle/libs.versions.toml @@ -33,6 +33,7 @@ lifecycleProcess = "2.8.7" agp = "8.7.2" kotlin = "1.9.25" livekit-uniffi = "0.1.12" +jna = "5.16.0" [libraries] livekit-uniffi = { module = "io.livekit:livekit-uniffi-android", version.ref = "livekit-uniffi" } @@ -103,6 +104,8 @@ mockito-inline = { module = "org.mockito:mockito-inline", version = "4.11.0" } byte-buddy = { module = "net.bytebuddy:byte-buddy", version = "1.14.3" } robolectric = { module = "org.robolectric:robolectric", version = "4.14.1" } +# JVM natives (libjnidispatch) for the Rust core under Robolectric; the AAR variant only ships Android ABIs. +jna = { module = "net.java.dev.jna:jna", version.ref = "jna" } turbine = { module = "app.cash.turbine:turbine", version = "1.0.0" } appcompat = { group = "androidx.appcompat", name = "appcompat", version.ref = "appcompat" } material = { group = "com.google.android.material", name = "material", version.ref = "material" } diff --git a/livekit-android-sdk/detekt-baseline-release.xml b/livekit-android-sdk/detekt-baseline-release.xml index 6bd3cb94..868b843c 100644 --- a/livekit-android-sdk/detekt-baseline-release.xml +++ b/livekit-android-sdk/detekt-baseline-release.xml @@ -21,7 +21,7 @@ CyclomaticComplexMethod:PreconnectAudioBuffer.kt$@Deprecated("Set AudioTrackPublishDefaults.preconnect = true on the RoomOptions instead.") suspend fun <T> Room.withPreconnectAudio( timeout: Duration = TIMEOUT, topic: String = DEFAULT_TOPIC, onError: ((e: Exception) -> Unit)? = null, operation: suspend () -> T, ) CyclomaticComplexMethod:PreconnectAudioBuffer.kt$internal suspend fun Room.startPreconnectAudioJob( roomScope: CoroutineScope, timeout: Duration = TIMEOUT, topic: String = DEFAULT_TOPIC ): () -> Unit CyclomaticComplexMethod:RTCEngine.kt$RTCEngine$@CheckResult internal suspend fun sendData(dataPacket: LivekitModels.DataPacket): Result<Unit> - CyclomaticComplexMethod:RTCEngine.kt$RTCEngine$@Synchronized @VisibleForTesting(otherwise = VisibleForTesting.PACKAGE_PRIVATE) fun reconnect() + CyclomaticComplexMethod:RTCEngine.kt$RTCEngine$@Synchronized internal fun reconnect(reason: ReconnectReason) CyclomaticComplexMethod:RTCEngine.kt$RTCEngine$fun onMessage(dataChannel: DataChannel, buffer: DataChannel.Buffer?) CyclomaticComplexMethod:RTCEngine.kt$RTCEngine$private fun makeRTCConfig( serverResponse: Either<JoinResponse, ReconnectResponse>, connectOptions: ConnectOptions, ): RTCConfiguration CyclomaticComplexMethod:Room.kt$Room$@Throws(Exception::class) suspend fun connect(url: String, token: String, options: ConnectOptions = ConnectOptions()) @@ -36,7 +36,7 @@ LargeClass:RTCEngine.kt$RTCEngine : Listener LargeClass:Room.kt$Room : ListenerParticipantListenerRpcManagerIncomingDataStreamManager LargeClass:SignalClient.kt$SignalClient : WebSocketListener - LongMethod:RTCEngine.kt$RTCEngine$@Synchronized @VisibleForTesting(otherwise = VisibleForTesting.PACKAGE_PRIVATE) fun reconnect() + LongMethod:RTCEngine.kt$RTCEngine$@Synchronized internal fun reconnect(reason: ReconnectReason) LongMethod:Room.kt$Room$@Throws(Exception::class) suspend fun connect(url: String, token: String, options: ConnectOptions = ConnectOptions()) LongMethod:SignalClient.kt$SignalClient$private fun handleSignalResponseImpl(ws: WebSocket, response: LivekitRtc.SignalResponse, encoded: ByteArray) LongParameterList:AudioBufferCallbackDispatcher.kt$AudioBufferCallback$(buffer: ByteBuffer, audioFormat: Int, channelCount: Int, sampleRate: Int, bytesRead: Int, captureTimeNs: Long) diff --git a/livekit-android-sdk/src/main/java/io/livekit/android/LiveKit.kt b/livekit-android-sdk/src/main/java/io/livekit/android/LiveKit.kt index c281b5aa..2a9a9076 100644 --- a/livekit-android-sdk/src/main/java/io/livekit/android/LiveKit.kt +++ b/livekit-android-sdk/src/main/java/io/livekit/android/LiveKit.kt @@ -23,6 +23,7 @@ import io.livekit.android.dagger.DaggerLiveKitComponent import io.livekit.android.dagger.RTCModule import io.livekit.android.dagger.create import io.livekit.android.room.Room +import io.livekit.android.telemetry.Telemetry import io.livekit.android.util.LKLog import io.livekit.android.util.LoggingLevel @@ -64,6 +65,17 @@ object LiveKit { @JvmStatic var enableWebRTCLogging: Boolean = false + /** + * Opts this process out of client telemetry, in effect when this returns. Collection stops, + * and everything not yet sent — queued, open or cached on disk — is deleted; Rooms created + * afterwards collect nothing, and the first of them deletes what a previous launch left cached. + * Call it at every launch, before creating a Room, to collect nothing at all. + * + * TODO: final shape pending the token/consent discussion. + */ + @JvmStatic + fun disableTelemetry() = Telemetry.disable() + /** * Certain WebRTC classes need to be initialized prior to use. * diff --git a/livekit-android-sdk/src/main/java/io/livekit/android/dagger/RTCModule.kt b/livekit-android-sdk/src/main/java/io/livekit/android/dagger/RTCModule.kt index bf4a692d..4288ea6a 100644 --- a/livekit-android-sdk/src/main/java/io/livekit/android/dagger/RTCModule.kt +++ b/livekit-android-sdk/src/main/java/io/livekit/android/dagger/RTCModule.kt @@ -38,6 +38,8 @@ import io.livekit.android.e2ee.DataPacketCryptorManagerImpl import io.livekit.android.memory.CloseableManager import io.livekit.android.room.datatrack.LocalDataTrackManagerFactory import io.livekit.android.room.datatrack.RemoteDataTrackManagerFactory +import io.livekit.android.telemetry.Telemetry +import io.livekit.android.telemetry.telemetryMicrophoneFailed import io.livekit.android.util.LKLog import io.livekit.android.util.LoggingLevel import io.livekit.android.webrtc.CustomAudioProcessingFactory @@ -113,6 +115,9 @@ internal object RTCModule { .setNativeLibraryName("lkjingle_peerconnection_so") .setInjectableLogger( { s, severity, s2 -> + if (severity == Logging.Severity.LS_ERROR) { + Telemetry.logWebRtc(s2, s) + } if (!LiveKit.enableWebRTCLogging) { return@setInjectableLogger } @@ -125,7 +130,8 @@ internal object RTCModule { else -> LoggingLevel.OFF } - LKLog.log(loggingLevel, null) { "$s2: $s" } + // The console only: telemetry already has WebRTC's errors, above. + if (loggingLevel >= LKLog.loggingLevel) LKLog.logger?.log(loggingLevel, null, "$s2: $s") }, Logging.Severity.LS_VERBOSE, ) @@ -182,6 +188,7 @@ internal object RTCModule { val audioRecordErrorCallback = object : JavaAudioDeviceModule.AudioRecordErrorCallback { override fun onWebRtcAudioRecordInitError(errorMessage: String?) { LKLog.e { "onWebRtcAudioRecordInitError: $errorMessage" } + telemetryMicrophoneFailed() } override fun onWebRtcAudioRecordStartError( @@ -189,10 +196,12 @@ internal object RTCModule { errorMessage: String?, ) { LKLog.e { "onWebRtcAudioRecordStartError: $errorCode. $errorMessage" } + telemetryMicrophoneFailed() } override fun onWebRtcAudioRecordError(errorMessage: String?) { LKLog.e { "onWebRtcAudioRecordError: $errorMessage" } + telemetryMicrophoneFailed() } } diff --git a/livekit-android-sdk/src/main/java/io/livekit/android/room/RTCEngine.kt b/livekit-android-sdk/src/main/java/io/livekit/android/room/RTCEngine.kt index c041ac66..5cf7ab4b 100644 --- a/livekit-android-sdk/src/main/java/io/livekit/android/room/RTCEngine.kt +++ b/livekit-android-sdk/src/main/java/io/livekit/android/room/RTCEngine.kt @@ -43,6 +43,9 @@ import io.livekit.android.room.util.MediaConstraintKeys import io.livekit.android.room.util.createAnswer import io.livekit.android.room.util.setLocalDescription import io.livekit.android.room.util.waitUntilConnected +import io.livekit.android.telemetry.RTCTelemetry +import io.livekit.android.telemetry.Telemetry +import io.livekit.android.telemetry.guarded import io.livekit.android.util.CloseableCoroutineScope import io.livekit.android.util.Either import io.livekit.android.util.FlowObservable @@ -65,9 +68,14 @@ import io.livekit.android.webrtc.peerconnection.RTCThreadToken import io.livekit.android.webrtc.peerconnection.executeBlockingOnRTCThread import io.livekit.android.webrtc.peerconnection.launchBlockingOnRTCThread import io.livekit.android.webrtc.toProtoSessionDescription +import io.livekit.uniffi.TelemetryScope +import io.livekit.uniffi.TelemetrySpan +import io.livekit.uniffi.telemetryDisconnectReason +import kotlinx.coroutines.CompletableDeferred import kotlinx.coroutines.CoroutineDispatcher import kotlinx.coroutines.Job import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.asContextElement import kotlinx.coroutines.coroutineScope import kotlinx.coroutines.delay import kotlinx.coroutines.ensureActive @@ -100,6 +108,10 @@ import livekit.org.webrtc.RtpSender import livekit.org.webrtc.RtpTransceiver import livekit.org.webrtc.RtpTransceiver.RtpTransceiverInit import livekit.org.webrtc.SessionDescription +import uniffi.livekit_telemetry.ReconnectReason +import uniffi.livekit_telemetry.SpanName +import uniffi.livekit_telemetry.SpanOutcome +import uniffi.livekit_telemetry.SpanStep import java.nio.ByteBuffer import javax.inject.Inject import javax.inject.Named @@ -109,6 +121,7 @@ import kotlin.coroutines.resume import kotlin.coroutines.resumeWithException import kotlin.time.Duration.Companion.milliseconds import kotlin.time.Duration.Companion.seconds +import uniffi.livekit_telemetry.DisconnectReason as FfiDisconnectReason /** * @suppress @@ -128,6 +141,37 @@ internal constructor( ) : SignalClient.Listener { internal var listener: Listener? = null + /** + * The Room's telemetry scope; null when telemetry is off. Bound on the engine's and the signal + * client's coroutines, so the Room handlers they drive log under the Room's session. + */ + internal var telemetryScope: TelemetryScope? = null + set(value) { + field = value + client.telemetryScope = value + } + + /** + * The Room's open `lk.connect` span while the user-initiated connect runs; the checkpoints + * are stamped here and in [SignalClient]. + */ + internal var connectSpan: TelemetrySpan? = null + set(value) { + field = value + client.connectSpan = value + } + + /** The Room's RTC instrument; the signal client hands it a manual subscribe's intent. */ + internal var rtcTelemetry: RTCTelemetry? + get() = client.rtcTelemetry + set(value) { + client.rtcTelemetry = value + } + + /** Why this session ended, for telemetry, when the SDK's enum says less: the server's Leave reason, or a reconnect that gave up. */ + @Volatile + internal var disconnectReasonForTelemetry: FfiDisconnectReason? = null + /** * When the current connection attempt began, taken at the top of [joinImpl]. Cleared once the * primary transport connects, so the attempt is timed exactly once. @@ -163,7 +207,7 @@ internal constructor( ConnectionState.DISCONNECTED -> { LKLog.d { "primary ICE disconnected" } if (oldVal == ConnectionState.CONNECTED) { - reconnect() + reconnect(if (isSubscriberPrimary) ReconnectReason.SUBSCRIBER_FAILED else ReconnectReason.PUBLISHER_FAILED) } } @@ -267,9 +311,11 @@ internal constructor( roomOptions: RoomOptions, ): JoinResponse { coroutineScope.close() - coroutineScope = CloseableCoroutineScope(SupervisorJob() + ioDispatcher) + coroutineScope = CloseableCoroutineScope(SupervisorJob() + ioDispatcher + Telemetry.currentScope.asContextElement(telemetryScope)) sessionUrl = url sessionToken = token + disconnectReasonForTelemetry = null + updateTelemetryServer() connectOptions = options lastRoomOptions = roomOptions return joinImpl(url, token, options, roomOptions) @@ -303,6 +349,10 @@ internal constructor( connectionState = ConnectionState.CONNECTING } val joinResponse = client.join(url, token, options, roomOptions) + guarded { + connectSpan?.step(SpanStep.Signal) + connectSpan?.step(SpanStep.JoinRecv) + } ensureActive() if (joinResponse.hasParticipant()) { @@ -325,6 +375,7 @@ internal constructor( isSubscriberPrimary = joinResponse.subscriberPrimary configure(joinResponse, options) + guarded { connectSpan?.step(SpanStep.PcCreated) } // The publisher created above needs the attempt's start time before its first offer, in // case video is published before the primary transport connects. publisher?.setConnectStartedAt(startedAtMs) @@ -398,7 +449,7 @@ internal constructor( // Also reconnect on publisher disconnect publisherObserver.connectionChangeListener = { newState -> if (newState.isDisconnected()) { - reconnect() + reconnect(ReconnectReason.PUBLISHER_FAILED) } } } else { @@ -605,9 +656,12 @@ internal constructor( /** * reconnect Signal and PeerConnections */ - @Synchronized @VisibleForTesting(otherwise = VisibleForTesting.PACKAGE_PRIVATE) - fun reconnect() { + fun reconnect() = reconnect(ReconnectReason.UNKNOWN) + + /** One reconnect cycle = one `lk.reconnect` span; attempts are its checkpoints. */ + @Synchronized + internal fun reconnect(reason: ReconnectReason) { if (reconnectingJob?.isActive == true) { LKLog.d { "Reconnection is already in progress" } return @@ -625,7 +679,8 @@ internal constructor( val forceFullReconnect = fullReconnectOnNext fullReconnectOnNext = false endSignalSession() - val job = coroutineScope.launch { + val reconnectSpan = guarded { telemetryScope?.start(SpanName.Reconnect(reason), null) } + val job = coroutineScope.launch(Telemetry.currentSpan.asContextElement(reconnectSpan)) { var hasResumedOnce = false var hasReconnectedOnce = false @@ -672,6 +727,7 @@ internal constructor( ReconnectType.FORCE_SOFT_RECONNECT -> false ReconnectType.FORCE_FULL_RECONNECT -> true } + guarded { reconnectSpan?.step(SpanStep.Attempt((retries + 1).toUInt(), isFullReconnect)) } var lastMessageSeq: Int? = null val connectOptions = connectOptions ?: ConnectOptions() @@ -784,6 +840,7 @@ internal constructor( outgoingDataTrackManager.republishTracks() } incomingDataTrackManager.resendSubscriptionUpdates() + guarded { reconnectSpan?.end(SpanOutcome.OK, null) } listener?.onPostReconnect(isFullReconnect) return@launch } @@ -795,12 +852,16 @@ internal constructor( } } + val gaveUp = !isClosed // else disconnect() won + guarded { if (gaveUp) reconnectSpan?.fail("ReconnectFailed") else reconnectSpan?.cancel() } + if (gaveUp) disconnectReasonForTelemetry = FfiDisconnectReason.RECONNECT_FAILED close("Failed reconnecting") listener?.onEngineDisconnected(DisconnectReason.UNKNOWN_REASON) } reconnectingJob = job job.invokeOnCompletion { + guarded { reconnectSpan?.takeIf { !it.isEnded() }?.cancel() } if (reconnectingJob == job) { reconnectingJob = null } @@ -1219,6 +1280,9 @@ internal constructor( internal const val TARGET_DATA_PACKET_SIZE = 15 * 1024 // 15 KB + /** A report WebRTC never delivers (its connection closed meanwhile) is skipped after this long. */ + private const val STATS_TIMEOUT_MS = 5_000L + /** * Corresponds to the max-message-size in SDP. Attempting to send packets * over this size will cause the data channel to close, so this must be enforced @@ -1370,7 +1434,7 @@ internal constructor( LKLog.i { "received close event: $reason, code: $code" } endSignalSession() abortPendingPublishTracks() - reconnect() + reconnect(ReconnectReason.SIGNAL_DISCONNECTED) } override fun onRemoteMuteChanged(trackSid: String, muted: Boolean) { @@ -1412,6 +1476,7 @@ internal constructor( else -> { close() + disconnectReasonForTelemetry = guarded { telemetryDisconnectReason(leave.reason.number) } val disconnectReason = leave.reason.convert() listener?.onEngineDisconnected(disconnectReason) } @@ -1444,6 +1509,14 @@ internal constructor( override fun onRefreshToken(token: String) { sessionToken = token regionUrlProvider?.token = token + updateTelemetryServer() + } + + /** Telemetry uploads with the Room's latest token, at join and on every refresh, to the URL the app gave. */ + private fun updateTelemetryServer() { + val url = regionUrlProvider?.serverUrl?.toString() ?: sessionUrl ?: return + val token = sessionToken ?: return + guarded { telemetryScope?.setServer(url, token) } } override fun onLocalTrackUnpublished(trackUnpublished: LivekitRtc.TrackUnpublishedResponse) { @@ -1663,6 +1736,24 @@ internal constructor( client.sendSyncState(syncState) } + /** Runs [action] on the RTC thread, suspending rather than blocking the caller. */ + internal suspend fun onRTCThread(action: () -> T): T? = launchBlockingOnRTCThread(rtcThreadToken) { action() } + + /** + * Each peer connection's whole report, publisher first, suspending rather than blocking. + * [request] makes each getStats() call, or refuses it (false): then no further one is made. + */ + internal suspend fun peerStats(request: (() -> Unit) -> Boolean): List { + val reports = mutableListOf() + for (transport in listOfNotNull(publisher, subscriber)) { + val report = CompletableDeferred() + val requested = transport.withPeerConnection { request { getStats { report.complete(it) } } } ?: continue + if (!requested) break + withTimeoutOrNull(STATS_TIMEOUT_MS) { report.await() }?.let(reports::add) + } + return reports + } + fun getPublisherRTCStats(callback: RTCStatsCollectorCallback) { runBlocking { publisher?.withPeerConnection { getStats(callback) } diff --git a/livekit-android-sdk/src/main/java/io/livekit/android/room/Room.kt b/livekit-android-sdk/src/main/java/io/livekit/android/room/Room.kt index 12e0121e..3dffcfb5 100644 --- a/livekit-android-sdk/src/main/java/io/livekit/android/room/Room.kt +++ b/livekit-android-sdk/src/main/java/io/livekit/android/room/Room.kt @@ -82,6 +82,12 @@ import io.livekit.android.room.track.Track import io.livekit.android.room.track.TrackPublication import io.livekit.android.room.types.toSDKType import io.livekit.android.room.util.ConnectionWarmer +import io.livekit.android.telemetry.RTCTelemetry +import io.livekit.android.telemetry.Telemetry +import io.livekit.android.telemetry.end +import io.livekit.android.telemetry.guarded +import io.livekit.android.telemetry.observeForTelemetry +import io.livekit.android.telemetry.telemetry import io.livekit.android.util.FlowObservable import io.livekit.android.util.LKLog import io.livekit.android.util.flow @@ -89,11 +95,14 @@ import io.livekit.android.util.flowDelegate import io.livekit.android.util.invoke import io.livekit.android.util.rethrowIfCancellationSignal import io.livekit.android.webrtc.getFilteredStats +import io.livekit.uniffi.TelemetryScope +import io.livekit.uniffi.TelemetrySpan import kotlinx.coroutines.CancellationException import kotlinx.coroutines.CoroutineDispatcher import kotlinx.coroutines.CoroutineExceptionHandler import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.asContextElement import kotlinx.coroutines.cancel import kotlinx.coroutines.coroutineScope import kotlinx.coroutines.ensureActive @@ -115,6 +124,11 @@ import livekit.org.webrtc.RendererCommon import livekit.org.webrtc.RtpReceiver import livekit.org.webrtc.SurfaceViewRenderer import livekit.org.webrtc.audio.AudioDeviceModule +import uniffi.livekit_telemetry.ReconnectReason +import uniffi.livekit_telemetry.RoomIdentity +import uniffi.livekit_telemetry.SpanName +import uniffi.livekit_telemetry.SpanOutcome +import uniffi.livekit_telemetry.SpanStep import java.net.URI import java.util.Date import javax.inject.Named @@ -168,9 +182,58 @@ constructor( private val eventBus = BroadcastEventBus() val events = eventBus.readOnly() + /** + * This Room's session in the telemetry pipeline, one trace for the Room's lifetime; null after + * [LiveKit.disableTelemetry][io.livekit.android.LiveKit.disableTelemetry]. + */ + @VisibleForTesting + internal var telemetryScope: TelemetryScope? = Telemetry.scope(context) + private val rtcTelemetry = telemetryScope?.let { RTCTelemetry(this, it) } + + /** Removes this Room's audio route and focus listeners: an app-supplied handler can outlive it. */ + private val stopAudioTelemetry = telemetryScope?.let { audioSwitchHandler?.observeForTelemetry() } + + /** The user-initiated connect, open from [connect] to [onEngineConnected]; the engine stamps its checkpoints. */ + private var connectSpan: TelemetrySpan? = null + set(value) { + field = value + engine.connectSpan = value + } + + /** + * Records an app event in this Room's telemetry, exported as `custom.` next to the SDK's + * own records, with this Room's correlation attributes. + * + * ``` + * room.emitTelemetryEvent("checkout.started", mapOf("cart.items" to "3")) + * ``` + * + * Names and keys up to 128 bytes, values up to 1024 bytes, at most 64 attributes and no `lk.` + * keys; anything else is dropped, never truncated. + */ + @JvmOverloads + fun emitTelemetryEvent(name: String, attributes: Map = emptyMap()) { + guarded { telemetryScope?.emitCustom(name, attributes) } + } + + /** + * Sets a correlation attribute on every telemetry record this Room captures from now on, to + * match them with your own data (an order id, a tenant); `null` removes it. + * + * ``` + * room.setTelemetryAttribute("app.order_id", order.id) + * ``` + * + * Same limits as [emitTelemetryEvent], at most 64 per Room. + */ + fun setTelemetryAttribute(key: String, value: String?) { + guarded { telemetryScope?.setAttribute(key, value) } + } + init { engine.listener = this - + engine.telemetryScope = telemetryScope + engine.rtcTelemetry = rtcTelemetry // Register SDK-internal text-stream handlers for the RPC v2 transport. These reserve // the topics `lk.rpc_request` and `lk.rpc_response` from user-level handler registration. incomingDataStreamManager.registerTextStreamHandler(RPC_REQUEST_DATA_STREAM_TOPIC) { receiver, fromIdentity -> @@ -365,6 +428,7 @@ constructor( */ val localParticipant: LocalParticipant = localParticipantFactory.create(dynacast = false).apply { internalListener = this@Room + telemetryScope = this@Room.telemetryScope } private var mutableRemoteParticipants by flowDelegate(emptyMap()) @@ -488,7 +552,14 @@ constructor( state = State.CONNECTING connectOptions = options - coroutineScope = CoroutineScope(defaultDispatcher + SupervisorJob()) + // The Room's destination and grant from the start, so an attempt failing before the join still uploads. + guarded { telemetryScope?.setServer(url, token) } + engine.disconnectReasonForTelemetry = null // this attempt's own, if it ends early + // One connect() = one attempt; reconnect cycles get their own spans. + connectSpan = guarded { telemetryScope?.start(SpanName.Connect, null) } + + coroutineScope = CoroutineScope(defaultDispatcher + SupervisorJob() + Telemetry.currentScope.asContextElement(telemetryScope)) + rtcTelemetry?.start(coroutineScope) roomOptions = getCurrentRoomOptions() @@ -513,7 +584,7 @@ constructor( // rethrow all throwables from the connect job. val emptyCoroutineExceptionHandler = CoroutineExceptionHandler { _, _ -> } val connectJob = coroutineScope.launch( - ioDispatcher + emptyCoroutineExceptionHandler, + ioDispatcher + emptyCoroutineExceptionHandler + Telemetry.currentSpan.asContextElement(connectSpan), ) { if (audioProcessingController is AuthedAudioProcessingController) { audioProcessingController.authenticate(url, token) @@ -609,6 +680,7 @@ constructor( connectJob.join() error?.let { + guarded { connectSpan?.end(it) } if (it !is CancellationException) { handleDisconnect(DisconnectReason.JOIN_FAILURE) } @@ -670,6 +742,7 @@ constructor( */ fun release() { disconnect() + stopAudioTelemetry?.invoke() closeableManager.close() } @@ -702,6 +775,7 @@ constructor( localParticipant.updateFromInfo(response.participant) localParticipant.setEnabledPublishCodecs(response.enabledPublishCodecsList) + updateTelemetryRoom() if (response.otherParticipantsList.isNotEmpty()) { response.otherParticipantsList.forEach { info -> @@ -710,6 +784,20 @@ constructor( } } + /** The room and local participant on every telemetry record of this session from now on. */ + private fun updateTelemetryRoom() { + guarded { + telemetryScope?.setRoom( + RoomIdentity( + sid = sid?.sid?.takeIf { it.isNotEmpty() }, + name = name, + participantSid = localParticipant.sid.value.takeIf { it.isNotEmpty() }, + participantIdentity = localParticipant.identity?.value, + ), + ) + } + } + private fun setupLocalParticipantEventHandling() { coroutineScope.launch { localParticipant.events.collect { @@ -1053,7 +1141,7 @@ constructor( if (state == State.RECONNECTING) { return } - engine.reconnect() + engine.reconnect(ReconnectReason.NETWORK_CHANGED) } private fun handleDisconnect(reason: DisconnectReason) { @@ -1069,6 +1157,12 @@ constructor( hasLostConnectivity = false state = State.DISCONNECTED + guarded { + connectSpan?.run { if (reason == DisconnectReason.CLIENT_INITIATED) cancel() else fail(reason.name) } + // Once per real session: never on a reconnect. + telemetryScope?.disconnected(engine.disconnectReasonForTelemetry ?: reason.telemetry) + } + connectSpan = null cleanupRoom() engine.close() @@ -1234,6 +1328,15 @@ constructor( */ override fun onEngineConnected() { state = State.CONNECTED + guarded { + connectSpan?.run { + step(SpanStep.Engine) + step(SpanStep.PcConnected) + step(SpanStep.RoomConnected) + end(SpanOutcome.OK, null) + } + } + connectSpan = null eventBus.postEvent(RoomEvent.Connected(this), coroutineScope) } @@ -1345,6 +1448,7 @@ constructor( override fun onRoomUpdate(update: LivekitModels.Room) { if (update.sid != null) { sid = Sid(update.sid) + updateTelemetryRoom() } val oldMetadata = metadata metadata = update.metadata diff --git a/livekit-android-sdk/src/main/java/io/livekit/android/room/SignalClient.kt b/livekit-android-sdk/src/main/java/io/livekit/android/room/SignalClient.kt index 788ca4da..5d14d364 100644 --- a/livekit-android-sdk/src/main/java/io/livekit/android/room/SignalClient.kt +++ b/livekit-android-sdk/src/main/java/io/livekit/android/room/SignalClient.kt @@ -27,6 +27,9 @@ import io.livekit.android.room.participant.ParticipantTrackPermission import io.livekit.android.room.track.Track import io.livekit.android.stats.NetworkInfo import io.livekit.android.stats.getClientInfo +import io.livekit.android.telemetry.RTCTelemetry +import io.livekit.android.telemetry.Telemetry +import io.livekit.android.telemetry.guarded import io.livekit.android.util.CloseableCoroutineScope import io.livekit.android.util.Either import io.livekit.android.util.LKLog @@ -36,12 +39,15 @@ import io.livekit.android.util.toHttpUrl import io.livekit.android.util.toWebsocketUrl import io.livekit.android.util.withDeadline import io.livekit.android.webrtc.toProtoSessionDescription +import io.livekit.uniffi.TelemetryScope +import io.livekit.uniffi.TelemetrySpan import kotlinx.coroutines.CancellableContinuation import kotlinx.coroutines.CompletableDeferred import kotlinx.coroutines.CoroutineDispatcher import kotlinx.coroutines.ExperimentalCoroutinesApi import kotlinx.coroutines.Job import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.asContextElement import kotlinx.coroutines.delay import kotlinx.coroutines.flow.MutableSharedFlow import kotlinx.coroutines.launch @@ -64,6 +70,7 @@ import okhttp3.WebSocket import okhttp3.WebSocketListener import okio.ByteString import okio.ByteString.Companion.toByteString +import uniffi.livekit_telemetry.SpanStep import java.util.Date import java.util.concurrent.ConcurrentHashMap import java.util.concurrent.atomic.AtomicInteger @@ -108,6 +115,15 @@ constructor( @Volatile private var joinContinuation: CancellableContinuation? = null + + /** The Room's open `lk.connect` span, for the signaling checkpoints; null outside the user-initiated connect. */ + internal var connectSpan: TelemetrySpan? = null + + /** The Room's telemetry scope: signal responses drive the Room's handlers, whose records are filed under it. */ + internal var telemetryScope: TelemetryScope? = null + + /** The Room's RTC instrument, for a manual subscribe's intent. */ + internal var rtcTelemetry: RTCTelemetry? = null private lateinit var coroutineScope: CloseableCoroutineScope /** @@ -197,7 +213,7 @@ constructor( LKLog.i { "connecting to $wsUrlString" } - coroutineScope = CloseableCoroutineScope(SupervisorJob() + ioDispatcher) + coroutineScope = CloseableCoroutineScope(SupervisorJob() + ioDispatcher + Telemetry.currentScope.asContextElement(telemetryScope)) lastUrl = wsUrlString lastOptions = options lastRoomOptions = roomOptions @@ -312,6 +328,10 @@ constructor( } // --------------------------------- WebSocket Listener --------------------------------------// + override fun onOpen(webSocket: WebSocket, response: Response) { + guarded { connectSpan?.step(SpanStep.WsOpen) } + } + override fun onMessage(webSocket: WebSocket, text: String) { if (webSocket != currentWs) { // Possibly message from old websocket, discard. @@ -439,6 +459,7 @@ constructor( } fun sendOffer(offer: SessionDescription, offerId: Int) { + guarded { connectSpan?.step(SpanStep.OfferSent) } val sd = offer.toProtoSessionDescription(offerId) val request = LivekitRtc.SignalRequest.newBuilder() .setOffer(sd) @@ -448,6 +469,7 @@ constructor( } fun sendAnswer(answer: SessionDescription, offerId: Int) { + guarded { connectSpan?.step(SpanStep.AnswerSent) } val sd = answer.toProtoSessionDescription(offerId) val request = LivekitRtc.SignalRequest.newBuilder() .setAnswer(sd) diff --git a/livekit-android-sdk/src/main/java/io/livekit/android/room/participant/LocalParticipant.kt b/livekit-android-sdk/src/main/java/io/livekit/android/room/participant/LocalParticipant.kt index 9ec1d98f..ff5165d3 100644 --- a/livekit-android-sdk/src/main/java/io/livekit/android/room/participant/LocalParticipant.kt +++ b/livekit-android-sdk/src/main/java/io/livekit/android/room/participant/LocalParticipant.kt @@ -65,16 +65,24 @@ import io.livekit.android.room.track.VideoPreset import io.livekit.android.room.track.screencapture.ScreenCaptureParams import io.livekit.android.room.util.EncodingUtils import io.livekit.android.rpc.RpcError +import io.livekit.android.telemetry.Telemetry +import io.livekit.android.telemetry.end +import io.livekit.android.telemetry.guarded +import io.livekit.android.telemetry.spanTrack import io.livekit.android.util.LKLog import io.livekit.android.util.flow import io.livekit.android.util.rethrowIfCancellationSignal import io.livekit.android.webrtc.sortVideoCodecPreferences +import io.livekit.uniffi.TelemetryScope import kotlinx.coroutines.CoroutineDispatcher import kotlinx.coroutines.Job import kotlinx.coroutines.NonCancellable +import kotlinx.coroutines.asContextElement import kotlinx.coroutines.async import kotlinx.coroutines.coroutineScope +import kotlinx.coroutines.currentCoroutineContext import kotlinx.coroutines.ensureActive +import kotlinx.coroutines.isActive import kotlinx.coroutines.launch import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.withLock @@ -95,6 +103,8 @@ import livekit.org.webrtc.RtpTransceiver.RtpTransceiverInit import livekit.org.webrtc.SurfaceTextureHelper import livekit.org.webrtc.VideoCapturer import livekit.org.webrtc.VideoProcessor +import uniffi.livekit_telemetry.SpanName +import uniffi.livekit_telemetry.SpanOutcome import java.nio.ByteBuffer import java.nio.charset.CodingErrorAction import java.util.Collections @@ -152,6 +162,9 @@ internal constructor( internal val enabledPublishVideoCodecs = Collections.synchronizedList(mutableListOf()) + /** The Room's telemetry scope, for the `lk.publish` span; null when telemetry is off. */ + internal var telemetryScope: TelemetryScope? = null + private var defaultAudioTrack: LocalAudioTrack? = null private var defaultVideoTrack: LocalVideoTrack? = null @@ -500,7 +513,7 @@ internal constructor( ) var publication: LocalTrackPublication? = null try { - publication = publishTrackImpl( + publication = publishTrackSpanned( track = track, options = options, requestConfig = { @@ -595,7 +608,7 @@ internal constructor( var publication: LocalTrackPublication? = null try { - publication = publishTrackImpl( + publication = publishTrackSpanned( track = track, options = options, requestConfig = { @@ -652,11 +665,15 @@ internal constructor( } /** + * One publish = one `lk.publish` span, under a still-running ambient span (the connect span + * for a pre-connect microphone), and ambient itself while [publishTrackImpl] runs, so the + * publish's own warnings point at it. + * * @throws TrackException.PublishException thrown when the publish fails. see [TrackException.PublishException.message] for details. * @return true if the track publish was successful. */ @Throws(TrackException.PublishException::class) - private suspend fun publishTrackImpl( + private suspend fun publishTrackSpanned( track: Track, options: TrackPublishOptions, requestConfig: AddTrackRequest.Builder.() -> Unit, @@ -667,8 +684,30 @@ internal constructor( LKLog.w { "Attempting to publish a disposed track, ignoring." } return null } + val parent = Telemetry.currentSpan.get()?.takeIf { guarded { !it.isEnded() } == true } + val span = guarded { telemetryScope?.start(SpanName.Publish, parent) } + try { + return withContext(Telemetry.currentSpan.asContextElement(span)) { + publishTrackImpl(track, options, requestConfig, encodings, publishListener) + } + } finally { + // A cancellation that wins before withContext runs the publish leaves the span to us. + val active = currentCoroutineContext().isActive + guarded { span?.takeIf { !it.isEnded() }?.run { if (active) fail("PublishException") else cancel() } } + } + } + @Throws(TrackException.PublishException::class) + private suspend fun publishTrackImpl( + track: Track, + options: TrackPublishOptions, + requestConfig: AddTrackRequest.Builder.() -> Unit, + encodings: List = emptyList(), + publishListener: PublishListener? = null, + ): LocalTrackPublication? { + val span = Telemetry.currentSpan.get() // this publish's own, from publishTrackSpanned fun onPublishFailure(e: TrackException.PublishException, triggerEvent: Boolean = true) { + guarded { span?.end(e) } publishListener?.onPublishFailure(e) if (triggerEvent) { eventBus.postEvent(ParticipantEvent.LocalTrackPublicationFailed(this, track, e), scope) @@ -680,6 +719,7 @@ internal constructor( } val trackSource = Track.Source.fromProto(addTrackRequestBuilder.source ?: LivekitModels.TrackSource.UNRECOGNIZED) + guarded { spanTrack(track.kind, trackSource)?.let { span?.setTrack(it) } } if (!hasPermissionsToPublish(trackSource)) { val exception = TrackException.PublishException("Failed to publish track, insufficient permissions") onPublishFailure(exception) @@ -856,6 +896,10 @@ internal constructor( participant = this, options = options, ) + guarded { + spanTrack(track.kind, trackSource, publication.sid)?.let { span?.setTrack(it) } + span?.end(SpanOutcome.OK, null) + } addTrackPublication(publication) LKLog.v { "add track publication $publication" } @@ -864,6 +908,8 @@ internal constructor( eventBus.postEvent(ParticipantEvent.LocalTrackPublished(this, publication), scope) } } finally { + val active = currentCoroutineContext().isActive + guarded { span?.takeIf { !it.isEnded() }?.run { if (active) fail("PublishException") else cancel() } } if (publication == null) { // Negotiation can win the race against a failed or cancelled add track request. // Without a publication there is no unpublish to stop the transceiver, so it diff --git a/livekit-android-sdk/src/main/java/io/livekit/android/room/track/RemoteTrackPublication.kt b/livekit-android-sdk/src/main/java/io/livekit/android/room/track/RemoteTrackPublication.kt index 0f04a588..0d572067 100644 --- a/livekit-android-sdk/src/main/java/io/livekit/android/room/track/RemoteTrackPublication.kt +++ b/livekit-android-sdk/src/main/java/io/livekit/android/room/track/RemoteTrackPublication.kt @@ -1,5 +1,5 @@ /* - * Copyright 2023-2025 LiveKit, Inc. + * Copyright 2023-2026 LiveKit, Inc. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -20,6 +20,7 @@ import io.livekit.android.dagger.InjectionNames import io.livekit.android.events.TrackEvent import io.livekit.android.events.collect import io.livekit.android.room.participant.RemoteParticipant +import io.livekit.android.telemetry.guarded import io.livekit.android.util.debounce import io.livekit.android.util.invoke import kotlinx.coroutines.CoroutineDispatcher @@ -132,6 +133,7 @@ class RemoteTrackPublication( build() } participant.signalClient.sendUpdateSubscription(isDesired, participantTracks) + if (subscribed) guarded { participant.signalClient.rtcTelemetry?.subscribeIntent(this, participant) } } /** diff --git a/livekit-android-sdk/src/main/java/io/livekit/android/room/track/video/CameraCapturerUtils.kt b/livekit-android-sdk/src/main/java/io/livekit/android/room/track/video/CameraCapturerUtils.kt index 2996e524..e88bdbb5 100644 --- a/livekit-android-sdk/src/main/java/io/livekit/android/room/track/video/CameraCapturerUtils.kt +++ b/livekit-android-sdk/src/main/java/io/livekit/android/room/track/video/CameraCapturerUtils.kt @@ -1,5 +1,5 @@ /* - * Copyright 2023-2025 LiveKit, Inc. + * Copyright 2023-2026 LiveKit, Inc. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -22,6 +22,7 @@ import android.content.Context import android.hardware.camera2.CameraManager import io.livekit.android.room.track.CameraPosition import io.livekit.android.room.track.LocalVideoTrackOptions +import io.livekit.android.telemetry.TelemetryCameraEvents import io.livekit.android.util.LKLog import livekit.org.webrtc.Camera1Capturer import livekit.org.webrtc.Camera1Enumerator @@ -98,6 +99,7 @@ object CameraCapturerUtils { ): Pair? { val cameraEnumerator = provider.provideEnumerator(context) val cameraEventsDispatchHandler = CameraEventsDispatchHandler() + cameraEventsDispatchHandler.registerHandler(TelemetryCameraEvents) val targetDevice = cameraEnumerator.findCamera(options.deviceId, options.position) ?: return null val targetVideoCapturer = provider.provideCapturer(context, options, cameraEventsDispatchHandler) diff --git a/livekit-android-sdk/src/main/java/io/livekit/android/telemetry/DeviceTelemetry.kt b/livekit-android-sdk/src/main/java/io/livekit/android/telemetry/DeviceTelemetry.kt new file mode 100644 index 00000000..a3d4f72c --- /dev/null +++ b/livekit-android-sdk/src/main/java/io/livekit/android/telemetry/DeviceTelemetry.kt @@ -0,0 +1,337 @@ +/* + * Copyright 2026 LiveKit, Inc. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package io.livekit.android.telemetry + +import android.app.Activity +import android.app.Application +import android.content.BroadcastReceiver +import android.content.ComponentCallbacks2 +import android.content.Context +import android.content.Intent +import android.content.IntentFilter +import android.content.res.Configuration +import android.media.AudioManager +import android.net.ConnectivityManager +import android.net.Network +import android.net.NetworkCapabilities +import android.os.BatteryManager +import android.os.Build +import android.os.Bundle +import android.os.PowerManager +import androidx.core.content.ContextCompat +import com.twilio.audioswitch.AudioDevice +import com.twilio.audioswitch.AudioDeviceChangeListener +import io.livekit.android.audio.AudioSwitchHandler +import io.livekit.android.util.LKLog +import io.livekit.uniffi.telemetrySetDeviceState +import kotlinx.coroutines.CoroutineExceptionHandler +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.ExperimentalCoroutinesApi +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.cancel +import kotlinx.coroutines.channels.Channel +import kotlinx.coroutines.launch +import livekit.org.webrtc.CameraVideoCapturer +import uniffi.livekit_telemetry.AppState +import uniffi.livekit_telemetry.AudioOutput +import uniffi.livekit_telemetry.AudioRouteReason +import uniffi.livekit_telemetry.CaptureDevice +import uniffi.livekit_telemetry.CaptureFailure +import uniffi.livekit_telemetry.DeviceEvent +import uniffi.livekit_telemetry.DeviceState +import uniffi.livekit_telemetry.MemoryPressure +import uniffi.livekit_telemetry.NetworkType +import uniffi.livekit_telemetry.TelemetryInstrument +import uniffi.livekit_telemetry.ThermalState +import java.util.concurrent.atomic.AtomicBoolean + +/** + * The device instrument: thermal status, battery saver, memory pressure, network, battery and app + * lifecycle as [DeviceState]. Every OS callback sends its change into one channel, drained in + * order by one coroutine on this instrument's own serial dispatcher; the core turns the state into + * `lk.device.*` records and upload holds. Nothing polls. + */ +@OptIn(ExperimentalCoroutinesApi::class) +internal class DeviceTelemetry(context: Context) : TelemetryInstrument { + private val app = context.applicationContext + private val power = app.getSystemService(Context.POWER_SERVICE) as? PowerManager + private val connectivity = app.getSystemService(Context.CONNECTIVITY_SERVICE) as? ConnectivityManager + + @Suppress("InjectDispatcher") // the instrument's own serial dispatcher, outside any Room + private val scope = CoroutineScope( + SupervisorJob() + Dispatchers.Default.limitedParallelism(1) + CoroutineExceptionHandler { _, e -> LKLog.w(e) { "Device telemetry stopped." } }, + ) + + // ponytail: unbounded, OS callbacks are a handful per minute at most + private val changes = Channel Unit>(Channel.UNLIMITED) + private var thermalListener: PowerManager.OnThermalStatusChangedListener? = null + + private fun post(change: DeviceState.() -> Unit = {}) { + changes.trySend(change) + } + + // MARK: - Lifecycle + + override fun start() { + scope.launch { + val state = DeviceState( + thermal = ThermalState.UNKNOWN, + lowPowerMode = null, + appState = AppState.FOREGROUND, + memory = MemoryPressure.NORMAL, + network = NetworkType.UNKNOWN, + ) + for (change in changes) { + // One failing change must not end the stream. + runCatching { + state.change() + // Read fresh on every change: these have no payload of their own. + state.thermal = thermal() + state.lowPowerMode = power?.isPowerSaveMode // unknown without a power service + state.networkConstrained = Build.VERSION.SDK_INT >= Build.VERSION_CODES.N && + connectivity?.restrictBackgroundStatus == ConnectivityManager.RESTRICT_BACKGROUND_STATUS_ENABLED + telemetrySetDeviceState(state.copy()) + }.onFailure { e -> LKLog.w(e) { "Device telemetry skipped a change." } } + } + } + val filter = IntentFilter().apply { + addAction(Intent.ACTION_BATTERY_CHANGED) + addAction(PowerManager.ACTION_POWER_SAVE_MODE_CHANGED) + if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.N) addAction(ConnectivityManager.ACTION_RESTRICT_BACKGROUND_CHANGED) + } + // System broadcasts only; the return value is the sticky battery intent: the first push. + val sticky = ContextCompat.registerReceiver(app, receiver, filter, ContextCompat.RECEIVER_NOT_EXPORTED) + post { sticky?.let { battery(it) } } + // ponytail: no default-network callback below API 24, so the network stays unknown there + if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.N) { + connectivity?.registerDefaultNetworkCallback(networkCallback) // delivers the current network at once + } + if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.Q) { + thermalListener = PowerManager.OnThermalStatusChangedListener { post() }.also { power?.addThermalStatusListener(it) } + } + app.registerComponentCallbacks(memoryCallbacks) + (app as? Application)?.registerActivityLifecycleCallbacks(activityCallbacks) + } + + /** + * Called by the core on the caller's thread (the opt-out's, often main) under its lifecycle + * lock: only unregisters and cancels, never waits, never calls back into telemetry, never throws. + */ + override fun stop() { + runCatching { app.unregisterReceiver(receiver) } + runCatching { connectivity?.unregisterNetworkCallback(networkCallback) } + if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.Q) runCatching { thermalListener?.let { power?.removeThermalStatusListener(it) } } + runCatching { app.unregisterComponentCallbacks(memoryCallbacks) } + runCatching { (app as? Application)?.unregisterActivityLifecycleCallbacks(activityCallbacks) } + scope.cancel() + } + + // MARK: - Power, battery, Data Saver + + private val receiver = object : BroadcastReceiver() { + override fun onReceive(context: Context, intent: Intent) { + post { if (intent.action == Intent.ACTION_BATTERY_CHANGED) battery(intent) } + } + } + + private fun DeviceState.battery(intent: Intent) { + val level = intent.getIntExtra(BatteryManager.EXTRA_LEVEL, -1) + val scale = intent.getIntExtra(BatteryManager.EXTRA_SCALE, 100) + val status = intent.getIntExtra(BatteryManager.EXTRA_STATUS, BatteryManager.BATTERY_STATUS_UNKNOWN) + batteryLevel = if (level >= 0 && scale > 0) (level * 100 / scale).toUInt() else null + batteryCharging = status == BatteryManager.BATTERY_STATUS_CHARGING || status == BatteryManager.BATTERY_STATUS_FULL + } + + private fun thermal(): ThermalState { + if (Build.VERSION.SDK_INT < Build.VERSION_CODES.Q) return ThermalState.UNKNOWN // no thermal API + return when (power?.currentThermalStatus ?: return ThermalState.UNKNOWN) { + PowerManager.THERMAL_STATUS_MODERATE -> ThermalState.FAIR + PowerManager.THERMAL_STATUS_SEVERE -> ThermalState.SERIOUS + PowerManager.THERMAL_STATUS_CRITICAL, PowerManager.THERMAL_STATUS_EMERGENCY, PowerManager.THERMAL_STATUS_SHUTDOWN -> ThermalState.CRITICAL + else -> ThermalState.NOMINAL // NONE, LIGHT + } + } + + // MARK: - Network + + private val networkCallback = object : ConnectivityManager.NetworkCallback() { + override fun onCapabilitiesChanged(network: Network, capabilities: NetworkCapabilities) = post { + this.network = capabilities.type + networkExpensive = !capabilities.hasCapability(NetworkCapabilities.NET_CAPABILITY_NOT_METERED) + } + + override fun onLost(network: Network) = post { + this.network = NetworkType.UNAVAILABLE + networkExpensive = false + } + } + + private val NetworkCapabilities.type: NetworkType + get() = when { + hasTransport(NetworkCapabilities.TRANSPORT_VPN) -> NetworkType.VPN + hasTransport(NetworkCapabilities.TRANSPORT_WIFI) -> NetworkType.WIFI + hasTransport(NetworkCapabilities.TRANSPORT_CELLULAR) -> NetworkType.CELL + hasTransport(NetworkCapabilities.TRANSPORT_ETHERNET) -> NetworkType.WIRED + hasTransport(NetworkCapabilities.TRANSPORT_BLUETOOTH) -> NetworkType.BLUETOOTH + else -> NetworkType.OTHER + } + + // MARK: - Memory and app state + + /** + * The trim levels the OS sends; hiding every activity is the move to the background. The OS + * never says pressure is over: it counts as normal again when an activity starts. + */ + private val memoryCallbacks = object : ComponentCallbacks2 { + override fun onTrimMemory(level: Int) = when (level) { + ComponentCallbacks2.TRIM_MEMORY_UI_HIDDEN -> post { appState = AppState.BACKGROUND } + ComponentCallbacks2.TRIM_MEMORY_RUNNING_CRITICAL, ComponentCallbacks2.TRIM_MEMORY_COMPLETE -> post { memory = MemoryPressure.CRITICAL } + ComponentCallbacks2.TRIM_MEMORY_RUNNING_MODERATE, + ComponentCallbacks2.TRIM_MEMORY_RUNNING_LOW, + ComponentCallbacks2.TRIM_MEMORY_BACKGROUND, + ComponentCallbacks2.TRIM_MEMORY_MODERATE, + -> post { memory = MemoryPressure.WARNING } + + else -> {} + } + + override fun onLowMemory() = post { memory = MemoryPressure.CRITICAL } + + override fun onConfigurationChanged(newConfig: Configuration) {} + } + + private val activityCallbacks = object : Application.ActivityLifecycleCallbacks { + override fun onActivityStarted(activity: Activity) = post { + appState = AppState.FOREGROUND + memory = MemoryPressure.NORMAL + } + + override fun onActivityCreated(activity: Activity, savedInstanceState: Bundle?) {} + + override fun onActivityResumed(activity: Activity) {} + + override fun onActivityPaused(activity: Activity) {} + + override fun onActivityStopped(activity: Activity) {} + + override fun onActivitySaveInstanceState(activity: Activity, outState: Bundle) {} + + override fun onActivityDestroyed(activity: Activity) {} + } +} + +// MARK: - Device events + +/** + * Audio route changes and focus loss: events, not state, they explain audio glitches. Process + * events, so one listener pair per handler however many Rooms share it (an app-supplied handler + * can be); the last Room to release it removes them. The audio switch names no reason for a + * change; a focus loss is an interruption that ends on the matching gain. Returns this Room's + * release, which is idempotent. + */ +internal fun AudioSwitchHandler.observeForTelemetry(): () -> Unit { + val observers = synchronized(audioObservers) { audioObservers.getOrPut(this) { AudioObservers() }.also { it.rooms++ } } + reconcile(observers) + val released = AtomicBoolean(false) + return { + if (released.compareAndSet(false, true)) { + synchronized(audioObservers) { observers.rooms-- } + reconcile(observers) + } + } +} + +/** + * Registers or unregisters [observers] until that matches whether any Room uses them. The + * handler is only called outside the registry lock: it dispatches under its own listener locks, + * and a callback may create or release a Room. One thread acts at a time; a thread that finds + * another acting leaves it to re-check when done. + */ +private fun AudioSwitchHandler.reconcile(observers: AudioObservers) { + while (true) { + val register = synchronized(audioObservers) { + val wanted = observers.rooms > 0 + if (observers.busy) return + if (observers.registered == wanted) { + if (!wanted && audioObservers[this] === observers) audioObservers.remove(this) + return + } + observers.busy = true + wanted + } + try { + if (register) { + registerAudioDeviceChangeListener(observers.route) + registerOnAudioFocusChangeListener(observers.focus) + } else { + unregisterAudioDeviceChangeListener(observers.route) + unregisterOnAudioFocusChangeListener(observers.focus) + } + } finally { + synchronized(audioObservers) { + observers.registered = register + observers.busy = false + } + } + } +} + +/** The listener pair of one handler and how many Rooms use it; holds no handler or Room. Guarded by the registry. */ +private class AudioObservers { + var rooms = 0 + var registered = false + var busy = false + val route = object : AudioDeviceChangeListener { + override fun invoke(devices: List, selected: AudioDevice?) = + Telemetry.deviceEvent(DeviceEvent.AudioRouteChanged(listOfNotNull(selected?.output), AudioRouteReason.UNKNOWN)) + } + val focus = AudioManager.OnAudioFocusChangeListener { change -> Telemetry.deviceEvent(DeviceEvent.AudioInterruption(began = change < 0)) } +} + +// ponytail: strong keys; an entry lives exactly as long as some Room still holds its handler +private val audioObservers = HashMap() + +private val AudioDevice.output: AudioOutput + get() = when (this) { + is AudioDevice.BluetoothHeadset -> AudioOutput.BLUETOOTH + is AudioDevice.WiredHeadset -> AudioOutput.WIRED_HEADSET + is AudioDevice.Earpiece -> AudioOutput.RECEIVER + is AudioDevice.Speakerphone -> AudioOutput.SPEAKER + else -> AudioOutput.OTHER + } + +/** Camera failures from WebRTC's capturer, on every camera track the SDK opens. */ +internal object TelemetryCameraEvents : CameraVideoCapturer.CameraEventsHandler { + override fun onCameraError(message: String?) = + Telemetry.deviceEvent(DeviceEvent.CaptureFailed(CaptureDevice.CAMERA, CaptureFailure.OTHER)) + + override fun onCameraDisconnected() = + Telemetry.deviceEvent(DeviceEvent.CaptureFailed(CaptureDevice.CAMERA, CaptureFailure.DISCONNECTED)) + + override fun onCameraFreezed(message: String?) {} + + override fun onCameraOpening(cameraName: String?) {} + + override fun onFirstFrameAvailable() {} + + override fun onCameraClosed() {} +} + +/** A microphone that failed to start, from WebRTC's audio device module. */ +internal fun telemetryMicrophoneFailed() = + Telemetry.deviceEvent(DeviceEvent.CaptureFailed(CaptureDevice.MICROPHONE, CaptureFailure.OTHER)) diff --git a/livekit-android-sdk/src/main/java/io/livekit/android/telemetry/RTCTelemetry.kt b/livekit-android-sdk/src/main/java/io/livekit/android/telemetry/RTCTelemetry.kt new file mode 100644 index 00000000..6bb82b9b --- /dev/null +++ b/livekit-android-sdk/src/main/java/io/livekit/android/telemetry/RTCTelemetry.kt @@ -0,0 +1,200 @@ +/* + * Copyright 2026 LiveKit, Inc. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package io.livekit.android.telemetry + +import androidx.annotation.VisibleForTesting +import io.livekit.android.events.RoomEvent +import io.livekit.android.room.Room +import io.livekit.android.room.participant.RemoteParticipant +import io.livekit.android.room.track.RemoteTrackPublication +import io.livekit.android.room.track.TrackPublication +import io.livekit.android.util.LKLog +import io.livekit.uniffi.TelemetryScope +import kotlinx.coroutines.CoroutineExceptionHandler +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.CoroutineStart +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.channels.Channel +import kotlinx.coroutines.channels.ReceiveChannel +import kotlinx.coroutines.coroutineScope +import kotlinx.coroutines.flow.first +import kotlinx.coroutines.flow.takeWhile +import kotlinx.coroutines.isActive +import kotlinx.coroutines.launch +import kotlinx.coroutines.withTimeoutOrNull +import livekit.org.webrtc.RTCStatsReport +import uniffi.livekit_telemetry.AttributeValue +import uniffi.livekit_telemetry.RtcStat +import java.math.BigInteger +import kotlin.coroutines.cancellation.CancellationException + +/** + * A Room's RTC instrument. Reports the remote tracks' lifecycle, from which the core runs the + * `lk.subscribe` span (intent → first media), and hands it one raw `getStats()` report per peer + * connection as often as it asks; the core maps every RTP stream to its track and windows it. + */ +internal class RTCTelemetry(private val room: Room, private val scope: TelemetryScope) { + /** Cuts the current wait short when a track appears: the core then asks for its 1 s pace. */ + private val wake = Channel(Channel.CONFLATED) + + /** + * Runs for one connection of the Room, in [connection]. Started before the Room joins, so the + * tracks published and subscribed during the join are seen too. Cancelled at the opt-out, + * which never waits for it; a getStats() request still on its way is refused by [Telemetry.ifCollecting]. + */ + fun start(connection: CoroutineScope) { + // Fail-open like the rest of telemetry: a failing core call is logged, never the app's crash. + val failOpen = CoroutineExceptionHandler { _, e -> LKLog.w(e) { "RTC telemetry stopped." } } + // Undispatched: subscribed to the Room's events before this returns. + val collector = connection.launch(failOpen, start = CoroutineStart.UNDISPATCHED) { + room.events.events.takeWhile { !Telemetry.disabled }.collect { event -> + runCatching { onEvent(event) }.onFailure { LKLog.w(it) { "RTC telemetry skipped ${event::class.simpleName}." } } + } + } + // On its own clock, not the Room's dispatcher: the core paces it. + + @Suppress("InjectDispatcher") + val poller = connection.launch(Dispatchers.Default + failOpen) { + pollStats(interval = { scope.statsPollIntervalMs().toLong() }, wake = wake) { + try { + recordPeerStats() + } catch (e: CancellationException) { + throw e + } catch (e: Exception) { + LKLog.w(e) { "RTC telemetry skipped a poll." } + } + } + } + connection.launch { + Telemetry.optedOut.first { it } + collector.cancel() + poller.cancel() + } + } + + private fun onEvent(event: RoomEvent) { + when (event) { + // With autoSubscribe the intent exists the moment the track is known. + is RoomEvent.TrackPublished -> { + if ((event.publication as? RemoteTrackPublication)?.isDesired == true) { + spanTrack(event.publication, event.participant as RemoteParticipant)?.let(scope::subscribeStarted) + } + wake.trySend(Unit) + } + + // A full reconnect announces the remote tracks again, in its own join. + is RoomEvent.Connected, is RoomEvent.Reconnected -> reconcileJoinedTracks() + + is RoomEvent.TrackSubscribed -> { + spanTrack(event.publication, event.participant)?.let(scope::subscribed) + wake.trySend(Unit) // a manual subscribe starts its wait here + } + + is RoomEvent.TrackSubscriptionFailed -> scope.subscribeFailed(event.sid, event.exception.errorType()) + is RoomEvent.TrackUnsubscribed -> scope.trackEnded(event.publications.sid) + is RoomEvent.TrackUnpublished -> scope.trackEnded(event.publication.sid) + else -> {} + } + } + + /** A manual subscribe: its intent opens `lk.subscribe`, and the poller takes up the core's faster pace at once. */ + fun subscribeIntent(publication: RemoteTrackPublication, participant: RemoteParticipant) { + spanTrack(publication, participant)?.let(scope::subscribeStarted) + wake.trySend(Unit) + } + + /** The join announces tracks without a TrackPublished event: their intent starts at connect. */ + private fun reconcileJoinedTracks() { + for (participant in room.remoteParticipants.values) { + val pending = participant.trackPublications.values.filter { (it as? RemoteTrackPublication)?.isDesired == true && it.track == null } + pending.forEach { publication -> spanTrack(publication, participant)?.let(scope::subscribeStarted) } + } + wake.trySend(Unit) + } + + /** + * One report per peer connection, with every track this Room sends or receives. Each getStats() + * call and each submit runs under the opt-out's lock, so none begins once [Telemetry.disable] + * has returned; waiting for an answer holds no lock. + */ + @VisibleForTesting + internal suspend fun recordPeerStats() { + // One RTC thread hop for every id: each track's own read then runs in place. + val tracks = room.engine.onRTCThread> { + buildMap { + for (participant in listOf(room.localParticipant) + room.remoteParticipants.values) { + for (publication in participant.trackPublications.values) { + publication.track?.withRTCTrack(null) { id() }?.let { put(it, publication.sid) } + } + } + } + } ?: return + for (report in room.engine.peerStats { request -> Telemetry.ifCollecting(request) != null }) { + if (report.statsMap.isEmpty()) continue + Telemetry.ifCollecting { scope.recordPeerStats(report.telemetryStats, tracks, report.telemetryTimestampNs) } ?: return + } + } +} + +/** + * Calls [poll] every [interval] ms, asking for the interval again after each wait. A [wake] + * re-reads the interval (a new track shortens it) but keeps the deadline, so wakes arriving faster + * than the interval never starve the polls. + */ +internal suspend fun pollStats( + interval: () -> Long, + wake: ReceiveChannel, + now: () -> Long = { System.nanoTime() / 1_000_000 }, + poll: suspend () -> Unit, +) = coroutineScope { + var polledAt = now() + while (isActive) { + val due = polledAt + interval() - now() + if (due > 0 && withTimeoutOrNull(due) { wake.receive() } != null) continue + polledAt = now() + poll() + } +} + +internal fun spanTrack(publication: TrackPublication, participant: RemoteParticipant) = + spanTrack(publication.kind, publication.source, publication.sid, participant.identity?.value) + +/** + * Every entry with its standard members as the core takes them, nested maps flattened with a + * dot (`qualityLimitationDurations.cpu`). No member names are known here. + */ +internal val RTCStatsReport.telemetryStats: List + get() = statsMap.values.map { stat -> + RtcStat(kind = stat.type, id = stat.id, members = buildMap { flatten(stat.members, "", this) }) + } + +internal val RTCStatsReport.telemetryTimestampNs: ULong + get() = (timestampUs.coerceAtLeast(0.0) * 1000).toULong() + +private fun flatten(values: Map<*, *>, prefix: String, into: MutableMap) { + for ((key, value) in values) { + val name = if (prefix.isEmpty()) key.toString() else "$prefix.$key" + when (value) { + is Boolean -> into[name] = AttributeValue.Bool(value) + is Int, is Long, is Short, is Byte, is BigInteger -> into[name] = AttributeValue.Int((value as Number).toLong()) + is Number -> into[name] = AttributeValue.Double(value.toDouble()) + is String -> into[name] = AttributeValue.Str(value) + is Map<*, *> -> flatten(value, name, into) + else -> {} // sequences carry nothing the core reads + } + } +} diff --git a/livekit-android-sdk/src/main/java/io/livekit/android/telemetry/Telemetry.kt b/livekit-android-sdk/src/main/java/io/livekit/android/telemetry/Telemetry.kt new file mode 100644 index 00000000..86a70555 --- /dev/null +++ b/livekit-android-sdk/src/main/java/io/livekit/android/telemetry/Telemetry.kt @@ -0,0 +1,285 @@ +/* + * Copyright 2026 LiveKit, Inc. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package io.livekit.android.telemetry + +import android.content.Context +import android.os.Build +import androidx.annotation.VisibleForTesting +import io.livekit.android.Version +import io.livekit.android.events.DisconnectReason +import io.livekit.android.room.track.Track +import io.livekit.android.util.LKLog +import io.livekit.android.util.LoggingLevel +import io.livekit.android.util.executeAsync +import io.livekit.uniffi.TelemetryScope +import io.livekit.uniffi.TelemetrySpan +import io.livekit.uniffi.telemetryConfigure +import io.livekit.uniffi.telemetryDeviceEvent +import io.livekit.uniffi.telemetryDisable +import io.livekit.uniffi.telemetryDisconnectReason +import io.livekit.uniffi.telemetryLog +import io.livekit.uniffi.telemetryScope +import kotlinx.coroutines.flow.MutableStateFlow +import livekit.LivekitModels +import okhttp3.OkHttpClient +import okhttp3.Request +import okhttp3.RequestBody.Companion.toRequestBody +import uniffi.livekit_telemetry.DeviceEvent +import uniffi.livekit_telemetry.ExportException +import uniffi.livekit_telemetry.ExportRequest +import uniffi.livekit_telemetry.ExportResponse +import uniffi.livekit_telemetry.LogRecord +import uniffi.livekit_telemetry.LogSource +import uniffi.livekit_telemetry.Sdk +import uniffi.livekit_telemetry.Severity +import uniffi.livekit_telemetry.SpanTrack +import uniffi.livekit_telemetry.TelemetryConfig +import uniffi.livekit_telemetry.TelemetryResource +import uniffi.livekit_telemetry.TelemetryTransport +import uniffi.livekit_telemetry.TrackKind +import uniffi.livekit_telemetry.TrackSource +import java.io.File +import java.io.IOException +import kotlin.coroutines.cancellation.CancellationException +import uniffi.livekit_telemetry.DisconnectReason as FfiDisconnectReason + +/** + * Client telemetry. The pipeline — destination, token, batching, retries, cache, holds, stats + * mapping, span state — lives in the Rust core, one per process; Android installs it with the + * first Room, feeds it OS signals and moves its bytes. Every tuning value is the core's default. + */ +@PublishedApi +internal object Telemetry { + /** Null until the first Room installs the pipeline; false when it could not start. */ + @Volatile + @VisibleForTesting + internal var installed: Boolean? = null + + /** Flips once, at the opt-out; a flow so a connected Room's RTC instrument stops right away. */ + val optedOut = MutableStateFlow(false) + + @VisibleForTesting + internal var disabled: Boolean + get() = optedOut.value + set(value) { + optedOut.value = value + } + + /** + * The span the current coroutine works inside: child spans nest under it and warn/error + * records point at it. Bound with `asContextElement` around connect, a reconnect cycle and + * publish. + */ + val currentSpan = ThreadLocal() + + /** + * The Room the current coroutine works for: bound on the Room's, the engine's and the signal + * client's coroutines, so a Room handler's warning lands in that Room's trace. + */ + val currentScope = ThreadLocal() + + /** A new Room's scope, installing the pipeline first if this is the first Room; null when off. */ + fun scope(context: Context): TelemetryScope? { + if (disabled) { + // A previous launch's cache: the core can only purge a pipeline it has, and none installs now. + runCatching { storageDirectory(context).deleteRecursively() } + return null + } + if (installed == null) { + synchronized(this) { + if (installed == null) configure(context) + } + } + return if (active) guarded { telemetryScope() } else null + } + + /** + * Install (or replace) the process pipeline; refused by the core after an opt-out. Fail-open, a + * missing native library included: the app runs without telemetry rather than not at all. + */ + fun configure(context: Context): Boolean { + installed = install(context.applicationContext) + return installed == true + } + + private fun install(context: Context): Boolean = try { + val sdk = TelemetryResource( + sdk = Sdk.ANDROID, + sdkVersion = Version.CLIENT_VERSION, + osName = "android", + osVersion = Build.VERSION.RELEASE ?: "", + deviceModel = "${Build.MANUFACTURER} ${Build.MODEL}".trim(), + ) + val config = TelemetryConfig(sdk = sdk, storageDir = storageDirectory(context).path) + telemetryConfigure(config, OkHttpTelemetryTransport(), listOf(DeviceTelemetry(context))) + true + } catch (e: Throwable) { + diagnose(e, "Telemetry could not start; running without it.") + false + } + + /** The on-disk batch cache: the app's cache directory, which the OS may clear and backups skip. */ + fun storageDirectory(context: Context) = File(context.cacheDir, "livekit-telemetry") + + /** + * The process opt-out, in effect when this returns: the core stops capturing and purges what + * it has not sent. [disabled] stays for [scope], which deletes a previous launch's cache + * without installing a pipeline (the core purges only one it has). + */ + fun disable() { + synchronized(this) { disabled = true } // after this, ifCollecting starts nothing + try { + telemetryDisable() + } catch (e: Throwable) { // a missing native library included: the opt-out never fails the app + diagnose(e, "The opt-out did not reach the core; Rooms created from now on still collect nothing.") + } + } + + fun deviceEvent(event: DeviceEvent) { + if (active) guarded { telemetryDeviceEvent(event) } + } + + /** + * Runs [collect] only while telemetry is on, under the lock [disable] sets the flag with: no + * collection begins once it has returned. [collect] must only start work, never wait for it. + */ + fun ifCollecting(collect: () -> T): T? = synchronized(this) { if (disabled) null else collect() } + + /** Installed and not opted out: the only state in which anything reaches the core. */ + private val active get() = installed == true && !disabled + + /** Whether [LKLog] hands records at [level] to telemetry, whatever the console level: the core's floor. */ + @PublishedApi + internal fun captures(level: LoggingLevel): Boolean = + active && level >= LoggingLevel.WARN && level != LoggingLevel.OFF + + /** + * An SDK warning or error, filed under the ambient span, else the ambient Room, else the + * process. Telemetry's own lines never feed back into the pipeline. + */ + @PublishedApi + internal fun log(level: LoggingLevel, t: Throwable?, message: String) { + if (!captures(level)) return + guarded { + val caller = Throwable("caller").stackTrace.firstOrNull { frame -> + !frame.className.startsWith(LKLog::class.java.name) && !frame.className.startsWith(Telemetry::class.java.name) + } + if (caller?.className?.startsWith(OWN_PACKAGE) == true) return + val span = currentSpan.get()?.takeIf { !it.isEnded() } + val record = LogRecord( + severity = if (level == LoggingLevel.WARN) Severity.WARN else Severity.ERROR, + source = LogSource.SDK, + body = listOfNotNull(message.takeIf { it.isNotEmpty() }, t?.toString()).joinToString(": "), + logger = caller?.className?.substringAfterLast('.')?.substringBefore('$'), + function = caller?.methodName, + file = caller?.fileName, + line = caller?.lineNumber?.takeIf { it > 0 }?.toUInt(), + spanId = span?.context()?.spanId, + ) + val scope = currentScope.get() + if (span == null && scope != null) scope.log(record) else telemetryLog(record) + } + } + + /** WebRTC's own errors, next to (never instead of) the app's console logger. */ + fun logWebRtc(tag: String, message: String) { + if (active) guarded { telemetryLog(LogRecord(Severity.ERROR, LogSource.WEB_RTC, message.trim(), logger = tag)) } + } + + private const val OWN_PACKAGE = "io.livekit.android.telemetry." +} + +// MARK: - Transport + +/** + * Moves the core's requests: status, headers and body go back untouched and the core decides + * what they mean. Only a missing answer is an error. Follows no redirect — the request carries + * the participant token — so a 3xx comes back as the answer. + */ +@Suppress("SwallowedException") // the core takes a reason, not a cause +internal class OkHttpTelemetryTransport( + private val client: OkHttpClient = OkHttpClient.Builder().followRedirects(false).followSslRedirects(false).build(), +) : TelemetryTransport { + override suspend fun send(request: ExportRequest): ExportResponse { + val httpRequest = try { + Request.Builder() + .url(request.url) + .post(request.body.toRequestBody()) + .apply { request.headers.forEach { (name, value) -> header(name, value) } } + .build() + } catch (e: IllegalArgumentException) { + throw ExportException.Rejected("invalid request: ${e.message}") + } + val response = try { + client.newCall(httpRequest).executeAsync() + } catch (e: IOException) { + throw ExportException.Retryable(e.toString(), null) + } + return response.use { + ExportResponse(it.code.toUShort(), it.headers.toMap(), it.body?.bytes() ?: ByteArray(0)) + } + } +} + +// MARK: - Shared vocabulary + +/** + * A telemetry-only core call, which must never become the caller's failure (a core panic surfaces + * as an exception): reported with [diagnose] and swallowed. + */ +internal inline fun guarded(call: () -> T): T? = try { + call() +} catch (e: Throwable) { + diagnose(e, "Telemetry call failed; ignored.") + null +} + +/** + * Best effort: a console line about telemetry's own failure, never the caller's failure, even when + * the app's logger throws. Not through [LKLog.log], which would feed it back into the core. + */ +internal fun diagnose(e: Throwable, message: String) { + runCatching { if (LoggingLevel.WARN >= LKLog.loggingLevel) LKLog.logger?.log(LoggingLevel.WARN, e, message) } +} + +/** End on an exception: `cancelled` for a cancellation, `error` otherwise. */ +internal fun TelemetrySpan.end(error: Throwable) { + if (error is CancellationException) cancel() else fail(error.errorType()) +} + +/** `error.type`: the exception's class name. */ +internal fun Throwable.errorType(): String = javaClass.simpleName.ifEmpty { javaClass.name } + +internal fun spanTrack(kind: Track.Kind, source: Track.Source, sid: String? = null, remoteIdentity: String? = null): SpanTrack? { + val trackKind = when (kind) { + Track.Kind.AUDIO -> TrackKind.AUDIO + Track.Kind.VIDEO -> TrackKind.VIDEO + Track.Kind.UNRECOGNIZED -> return null + } + val trackSource = when (source) { + Track.Source.CAMERA -> TrackSource.CAMERA + Track.Source.MICROPHONE -> TrackSource.MICROPHONE + Track.Source.SCREEN_SHARE -> TrackSource.SCREEN_SHARE + Track.Source.SCREEN_SHARE_AUDIO -> TrackSource.SCREEN_SHARE_AUDIO + Track.Source.UNKNOWN -> TrackSource.UNKNOWN + } + return SpanTrack(sid, trackKind, trackSource, remoteIdentity) +} + +/** The shared enum by the protocol's number; the SDK enum mirrors the protocol's names. */ +internal val DisconnectReason.telemetry: FfiDisconnectReason + get() = telemetryDisconnectReason(runCatching { LivekitModels.DisconnectReason.valueOf(name).number }.getOrDefault(0)) diff --git a/livekit-android-sdk/src/main/java/io/livekit/android/util/LKLog.kt b/livekit-android-sdk/src/main/java/io/livekit/android/util/LKLog.kt index 9fa634f9..27908af4 100644 --- a/livekit-android-sdk/src/main/java/io/livekit/android/util/LKLog.kt +++ b/livekit-android-sdk/src/main/java/io/livekit/android/util/LKLog.kt @@ -16,6 +16,7 @@ package io.livekit.android.util +import io.livekit.android.telemetry.Telemetry import io.livekit.android.util.LoggingLevel.DEBUG import io.livekit.android.util.LoggingLevel.ERROR import io.livekit.android.util.LoggingLevel.INFO @@ -109,8 +110,14 @@ class LKLog { /** @suppress */ inline fun log(loggingLevel: LoggingLevel, t: Throwable? = null, crossinline message: (() -> String)) { - if (loggingLevel >= LKLog.loggingLevel) { - logger?.log(loggingLevel, t, message()) + val console = loggingLevel >= LKLog.loggingLevel + // Telemetry captures warnings and errors whatever the console level. + if (console || Telemetry.captures(loggingLevel)) { + val text = message() + if (console) { + logger?.log(loggingLevel, t, text) + } + Telemetry.log(loggingLevel, t, text) } } } diff --git a/livekit-android-test/build.gradle b/livekit-android-test/build.gradle index 749b42ce..61280206 100644 --- a/livekit-android-test/build.gradle +++ b/livekit-android-test/build.gradle @@ -35,6 +35,17 @@ android { testOptions { unitTests { includeAndroidResources = true + all { test -> + // Robolectric runs the real Rust core (livekit_uniffi) through JNA; point it at a host + // build of the library: -PlivekitUniffiLibraryPath=… or LIVEKIT_UNIFFI_LIBRARY_PATH. + def uniffiLibraryPath = project.findProperty('livekitUniffiLibraryPath') ?: System.getenv('LIVEKIT_UNIFFI_LIBRARY_PATH') + if (uniffiLibraryPath) { + test.systemProperty 'jna.library.path', uniffiLibraryPath + } + // UniFFI objects register with android.system.SystemCleaner (API 34+), which Robolectric + // implements on top of a JDK-internal cleaner. + test.jvmArgs '--add-exports=java.base/jdk.internal.ref=ALL-UNNAMED' + } } } lint { @@ -129,6 +140,7 @@ dependencies { testImplementation libs.junit testImplementation libs.robolectric + testImplementation libs.jna testImplementation libs.okhttp.mockwebserver testImplementation "org.jetbrains.kotlin:kotlin-reflect:$kotlin_version" kaptTest libs.dagger.compiler diff --git a/livekit-android-test/src/main/java/io/livekit/android/test/mock/MockPeerConnection.kt b/livekit-android-test/src/main/java/io/livekit/android/test/mock/MockPeerConnection.kt index 6a51cfeb..cd590b4c 100644 --- a/livekit-android-test/src/main/java/io/livekit/android/test/mock/MockPeerConnection.kt +++ b/livekit-android-test/src/main/java/io/livekit/android/test/mock/MockPeerConnection.kt @@ -205,8 +205,14 @@ class MockPeerConnection( return true } + /** What [getStats] delivers, how often it was asked, and how it answers (a test can hold the answer back). */ + var statsReport = RTCStatsReport(0, emptyMap()) + var statsRequests = 0 + var statsAnswer: (RTCStatsCollectorCallback) -> Unit = { it.onStatsDelivered(statsReport) } + override fun getStats(callback: RTCStatsCollectorCallback?) { - callback?.onStatsDelivered(RTCStatsReport(0, emptyMap())) + statsRequests++ + callback?.let(statsAnswer) } override fun setBitrate(min: Int?, current: Int?, max: Int?): Boolean { diff --git a/livekit-android-test/src/test/java/io/livekit/android/room/RoomTest.kt b/livekit-android-test/src/test/java/io/livekit/android/room/RoomTest.kt index cde24d5f..77e3da26 100644 --- a/livekit-android-test/src/test/java/io/livekit/android/room/RoomTest.kt +++ b/livekit-android-test/src/test/java/io/livekit/android/room/RoomTest.kt @@ -71,6 +71,7 @@ import org.mockito.kotlin.doSuspendableAnswer import org.mockito.kotlin.stub import org.mockito.kotlin.whenever import org.robolectric.RobolectricTestRunner +import uniffi.livekit_telemetry.ReconnectReason @ExperimentalCoroutinesApi @RunWith(RobolectricTestRunner::class) @@ -104,6 +105,8 @@ class RoomTest { override fun create(dynacast: Boolean): LocalParticipant { return Mockito.mock(LocalParticipant::class.java) .apply { + // A real participant's sid is empty until the join; the mock's value class getter would be null. + doReturn("").whenever(this).sid whenever(this.events).thenReturn( object : EventListenable { override val events: SharedFlow = MutableSharedFlow() @@ -229,7 +232,7 @@ class RoomTest { callback.onAvailable(network) } - Mockito.verify(rtcEngine).reconnect() + Mockito.verify(rtcEngine).reconnect(ReconnectReason.NETWORK_CHANGED) } @Test diff --git a/livekit-android-test/src/test/java/io/livekit/android/telemetry/TelemetryMockE2ETest.kt b/livekit-android-test/src/test/java/io/livekit/android/telemetry/TelemetryMockE2ETest.kt new file mode 100644 index 00000000..23275596 --- /dev/null +++ b/livekit-android-test/src/test/java/io/livekit/android/telemetry/TelemetryMockE2ETest.kt @@ -0,0 +1,362 @@ +/* + * Copyright 2026 LiveKit, Inc. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package io.livekit.android.telemetry + +import android.app.Application +import android.content.ComponentCallbacks2 +import android.content.Intent +import android.os.BatteryManager +import androidx.test.core.app.ApplicationProvider +import io.livekit.android.LiveKit +import io.livekit.android.room.ReconnectType +import io.livekit.android.room.SignalClient +import io.livekit.android.test.MockE2ETest +import io.livekit.android.test.mock.MockAudioStreamTrack +import io.livekit.android.test.mock.MockMediaStream +import io.livekit.android.test.mock.MockRtpReceiver +import io.livekit.android.test.mock.TestData +import io.livekit.android.test.mock.createMediaStreamId +import io.livekit.android.test.mock.room.track.createMockLocalAudioTrack +import io.livekit.android.util.LKLog +import io.livekit.uniffi.telemetryDiagnostics +import io.livekit.uniffi.telemetryFlush +import io.livekit.uniffi.telemetryScope +import io.livekit.uniffi.telemetryStats +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.ExperimentalCoroutinesApi +import kotlinx.coroutines.launch +import kotlinx.coroutines.withContext +import kotlinx.serialization.json.Json +import kotlinx.serialization.json.JsonElement +import kotlinx.serialization.json.JsonObject +import kotlinx.serialization.json.JsonPrimitive +import kotlinx.serialization.json.booleanOrNull +import kotlinx.serialization.json.contentOrNull +import kotlinx.serialization.json.intOrNull +import kotlinx.serialization.json.jsonArray +import kotlinx.serialization.json.jsonObject +import kotlinx.serialization.json.jsonPrimitive +import livekit.LivekitRtc +import livekit.org.webrtc.RTCStats +import livekit.org.webrtc.RTCStatsReport +import org.junit.Assert.assertEquals +import org.junit.Assert.assertFalse +import org.junit.Assert.assertNull +import org.junit.Assert.assertTrue +import org.junit.Assume.assumeTrue +import org.junit.Rule +import org.junit.Test +import org.junit.rules.TestRule +import org.junit.runner.RunWith +import org.junit.runners.model.Statement +import org.robolectric.RobolectricTestRunner +import uniffi.livekit_telemetry.AudioOutput +import uniffi.livekit_telemetry.AudioRouteReason +import uniffi.livekit_telemetry.DeviceEvent +import java.io.File +import java.math.BigInteger +import java.net.InetSocketAddress +import java.net.Socket +import java.util.UUID + +/** + * End to end through the Rust core into a local OpenTelemetry collector, on the SDK's mocks (mock + * websocket, mock peer connections): `otelcol-contrib --config src/test/resources/telemetry/otelcol.yaml` + * writes every OTLP request as a JSON line to [COLLECTOR_OUTPUT], `LK_TELEMETRY_ENDPOINT=http://127.0.0.1:4319` + * points the core at it, and a host build of `livekit_uniffi` runs it (`-PlivekitUniffiLibraryPath=…`). + * Skipped when any of them is missing. The pipeline is process-wide, hence one story. + */ +@ExperimentalCoroutinesApi +@RunWith(RobolectricTestRunner::class) +class TelemetryMockE2ETest : MockE2ETest() { + /** Unix nanoseconds when the pipeline started: the device instrument reports its initial state then. */ + private var startNs = 0L + + /** The pipeline is (re)installed before [mocksSetup] creates the Room, as the first Room of an app would. */ + @get:Rule + val telemetryRule = TestRule { base, _ -> + object : Statement() { + override fun evaluate() { + assumeTrue("LK_TELEMETRY_ENDPOINT points the core at a local collector", System.getenv("LK_TELEMETRY_ENDPOINT") != null) + assumeTrue("a collector listens on :4319", runCatching { Socket().use { it.connect(InetSocketAddress("127.0.0.1", 4319), 500) } }.isSuccess) + startNs = System.currentTimeMillis() * 1_000_000 + val app = ApplicationProvider.getApplicationContext() + // The OS's sticky battery broadcast, which the device instrument reads when it starts. + app.sendStickyBroadcast(Intent(Intent.ACTION_BATTERY_CHANGED).putExtra(BatteryManager.EXTRA_LEVEL, 42).putExtra(BatteryManager.EXTRA_SCALE, 100)) + assumeTrue("the pipeline starts (is livekit_uniffi on jna.library.path?)", Telemetry.configure(app)) + base.evaluate() + } + } + } + + /** + * One call and its aftermath: device changes → connect (a remote microphone announced by the + * join) → publish → subscribe to first media → app data and warnings → quick reconnect → + * disconnect → opt-out. + */ + @Test + fun aCallFromConnectToOptOut() = runTest { + val marker = UUID.randomUUID().toString() + val traceId = checkNotNull(room.telemetryScope).traceId() + room.setReconnectionType(ReconnectType.FORCE_SOFT_RECONNECT) + postDeviceChanges() + + connect( + TestData.JOIN.toBuilder() + .apply { join = join.toBuilder().addOtherParticipants(TestData.REMOTE_PARTICIPANT).build() } + .build(), + ) + + // App data: a correlation attribute on everything from now on (one set, one removed). + room.setTelemetryAttribute("app.call_id", marker) + room.setTelemetryAttribute("app.removed", marker) + room.setTelemetryAttribute("app.removed", null) + + // Publish a mock microphone: the server's TrackPublished answer completes the request. + val publish = launch { room.localParticipant.publishAudioTrack(createMockLocalAudioTrack()) } + simulateMessageFromServer(TestData.LOCAL_TRACK_PUBLISHED) + publish.join() + + // Subscribe: the announced remote microphone's track arrives, then its media, and ours leaves. + room.onAddTrack( + MockRtpReceiver.create(), + MockAudioStreamTrack(id = REMOTE_TRACK_ID), + arrayOf(MockMediaStream(id = createMediaStreamId(TestData.REMOTE_PARTICIPANT.sid, TestData.REMOTE_AUDIO_TRACK.sid))), + ) + getSubscriberPeerConnection().statsReport = report( + RTCStats(0, "inbound-rtp", "IN", mapOf("kind" to "audio", "trackIdentifier" to REMOTE_TRACK_ID, "bytesReceived" to BigInteger("1200"), "packetsReceived" to 10L)), + ) + getPublisherPeerConnection().statsReport = report( + RTCStats(0, "media-source", "MS", mapOf("kind" to "audio", "trackIdentifier" to TestData.LOCAL_TRACK_PUBLISHED.trackPublished.cid)), + RTCStats(0, "outbound-rtp", "OUT", mapOf("kind" to "audio", "mediaSourceId" to "MS", "bytesSent" to BigInteger("800"), "packetsSent" to 8L)), + ) + Thread.sleep(3000) // first media, at the core's 1 s polls + + // A Room handler's warning with no span in flight lands in the Room's session: the server + // unpublishes a track we never had. Outside any Room it is the process's. + val unknownSid = "TR_unknown_$marker" + simulateMessageFromServer( + LivekitRtc.SignalResponse.newBuilder() + .setTrackUnpublished(LivekitRtc.TrackUnpublishedResponse.newBuilder().setTrackSid(unknownSid)) + .build(), + ) + LKLog.e { "$marker process" } + room.emitTelemetryEvent("e2e.checkpoint", mapOf("e2e.marker" to marker)) + + // A quick reconnect on the mocks: the primary (subscriber) peer connection fails, the + // websocket reconnects, ICE reconnects. + disconnectPeerConnection() + testScheduler.advanceTimeBy(1000) + reconnectWebsocket() + connectPeerConnection() + testScheduler.advanceTimeBy(1000) + Thread.sleep(1500) // a poll on the reconnected session + + room.disconnect() + // The device state is pushed from the instrument's own dispatcher: wait until it has shipped. + val otlp = flushed { file -> file.logs.any { it.eventName == "lk.device.thermal.changed" } } + // The pipeline's own account, for a run whose records did not all arrive. + val diagnostics = telemetryDiagnostics() + println("telemetry: $diagnostics") + val spans = otlp.spans.filter { it.traceId == traceId } + val logs = otlp.logs.filter { it.traceId == traceId } + + // Connect: one span with the required checkpoints; the reconnect is its own span. + val connect = spans.single("lk.connect") + assertTrue(connect.events.toString(), connect.events.containsAll(listOf("ws_open", "signal", "join_recv", "pc_created", "engine", "pc_connected", "room_connected"))) + assertEquals("ok", connect.attributes["lk.outcome"]) + assertEquals("1", connect.attributes["lk.connect.attempt"]) + val reconnect = spans.single("lk.reconnect") + assertEquals(reconnect.attributes.toString(), "subscriber_failed", reconnect.attributes["lk.reconnect.reason"]) + assertEquals("ok", reconnect.attributes["lk.outcome"]) + assertTrue(reconnect.events.toString(), "attempt 1 quick" in reconnect.events) + + // Publish, and the subscribe of a track the join announced: intent → subscribed → first media. + val published = spans.single("lk.publish") + assertEquals("ok", published.attributes["lk.outcome"]) + assertEquals(TestData.LOCAL_AUDIO_TRACK.sid, published.attributes["lk.track.sid"]) + assertEquals("microphone", published.attributes["lk.track.source"]) + val subscribes = spans.filter { it.name == "lk.subscribe" }.associateBy { it.attributes["lk.track.sid"] } + val audio = checkNotNull(subscribes[TestData.REMOTE_AUDIO_TRACK.sid]) { "${subscribes.keys}" } + assertEquals(audio.attributes.toString(), "ok", audio.attributes["lk.outcome"]) + assertEquals(TestData.REMOTE_PARTICIPANT.identity, audio.attributes["lk.participant.remote_identity"]) + assertEquals(listOf("subscribed", "first_media"), audio.events) + // The announced camera never arrives: its intent is a span too, ended by the disconnect. + val video = checkNotNull(subscribes[TestData.REMOTE_VIDEO_TRACK.sid]) { "${subscribes.keys}" } + assertEquals(video.attributes.toString(), "cancelled", video.attributes["lk.outcome"]) + + // RTC windows from one report per peer connection, each track in its direction. + val windows = logs.filter { it.eventName == "lk.rtc.stats.sample" } + for ((sid, direction) in listOf(TestData.LOCAL_AUDIO_TRACK.sid to "outbound", TestData.REMOTE_AUDIO_TRACK.sid to "inbound")) { + val window = windows.any { it.attributes["lk.track.sid"] == sid && it.attributes["lk.track.direction"] == direction } + assertTrue("$direction window: ${windows.map { it.attributes }}", window) + } + assertTrue("windows carry the correlation attribute", windows.all { it.attributes["app.call_id"] == marker }) + + // App data and SDK records. + assertTrue(logs.any { it.eventName == "custom.e2e.checkpoint" && it.attributes["e2e.marker"] == marker && it.attributes["app.call_id"] == marker }) + assertFalse(otlp.logs.any { it.attributes["app.removed"] != null }) + val handler = logs.singleOrNull { it.body?.endsWith(unknownSid) == true } + assertTrue("a Room handler's warning: the Room's session, no span: $handler", handler != null && handler.spanId.isEmpty()) + assertEquals("roomname", handler!!.attributes["lk.room.name"]) + assertEquals(TestData.LOCAL_PARTICIPANT.identity, handler.attributes["lk.participant.identity"]) + val process = otlp.logs.singleOrNull { it.body == "$marker process" } + assertTrue("outside any Room: the process scope", process != null && process.traceId != traceId) + assertTrue("log records are warnings and errors only", otlp.logs.filter { it.eventName.isEmpty() }.all { it.severity >= SEVERITY_WARN }) + + // The session ends once, never on a reconnect. + val ended = logs.filter { it.eventName == "lk.room.disconnected" } + assertEquals(ended.map { it.attributes }.toString(), listOf("client_initiated"), ended.map { it.attributes["lk.disconnect.reason"] }) + + // Device: the state's initial values, what Robolectric can drive, and what it can only post. + val device = otlp.logs.filter { it.eventName.startsWith("lk.device.") } + for (event in listOf("thermal", "low_power", "network", "battery", "memory", "app_state").map { "lk.device.$it.changed" } + DEVICE_EVENTS) { + assertTrue("$event: ${device.map { it.eventName }} ($diagnostics)", device.any { it.eventName == event }) + } + assertTrue(device.any { it.eventName == "lk.device.memory.changed" && it.attributes.containsValue("critical") }) + assertTrue(device.any { it.eventName == "lk.device.app_state.changed" && it.attributes.containsValue("background") }) + val captures = device.filter { it.eventName == "lk.device.capture.failed" }.map { it.attributes["lk.device.capture.device"] to it.attributes["lk.device.capture.reason"] } + assertTrue("$captures", captures.containsAll(listOf("camera" to "disconnected", "microphone" to "other"))) + val stats = checkNotNull(telemetryStats()) + assertEquals("the whole call shipped", 0uL, stats.dropped) + + expectOptOut(marker) + } + + private suspend fun expectOptOut(marker: String) { + // Opt-out: what was not yet sent is deleted, nothing is collected afterwards. + val pending = component.roomFactory().create(context) + pending.emitTelemetryEvent("$marker.pending") + LiveKit.disableTelemetry() + assertNull("the opt-out is in effect when it returns", telemetryScope()) + assertNull("a Room created after the opt-out collects nothing", component.roomFactory().create(context).telemetryScope) + assertFalse("an unsent event is deleted, never uploaded", flushed().logs.any { it.eventName == "custom.$marker.pending" }) + assertEquals("the on-disk cache is purged", emptyList(), Telemetry.storageDirectory(context).list()?.toList().orEmpty()) + } + + /** Ships what the core holds and reads back what the collector wrote since the pipeline started. */ + private suspend fun flushed(until: (OtlpFile) -> Boolean = { true }): OtlpFile { + var file = OtlpFile(File(COLLECTOR_OUTPUT), since = startNs) + for (attempt in 0 until 5) { + withContext(Dispatchers.IO) { + // Each flush uploads a few batches: repeat until none is left. + for (i in 0 until 10) { + telemetryFlush() + if ((telemetryStats()?.cachedBatches ?: 0uL) == 0uL) break + } + } + Thread.sleep(2000) // the collector's file write + file = OtlpFile(File(COLLECTOR_OUTPUT), since = startNs) + if (until(file)) break + } + return file + } + + /** OS changes Robolectric can drive through the real callbacks, and events posted where it cannot. */ + private fun postDeviceChanges() { + val app = ApplicationProvider.getApplicationContext() + app.onTrimMemory(ComponentCallbacks2.TRIM_MEMORY_RUNNING_CRITICAL) + app.onTrimMemory(ComponentCallbacks2.TRIM_MEMORY_UI_HIDDEN) + TelemetryCameraEvents.onCameraDisconnected() + telemetryMicrophoneFailed() + // No audio switch in the mocks (NoAudioHandler): posted the way its listeners would. + Telemetry.deviceEvent(DeviceEvent.AudioRouteChanged(listOf(AudioOutput.SPEAKER), AudioRouteReason.UNKNOWN)) + Telemetry.deviceEvent(DeviceEvent.AudioInterruption(began = true)) + Thread.sleep(500) + } + + private fun reconnectWebsocket() { + wsFactory.listener.onOpen(wsFactory.ws, createOpenResponse(wsFactory.request)) + val softReconnectParam = wsFactory.request.url.queryParameter(SignalClient.CONNECT_QUERY_RECONNECT)?.toIntOrNull() ?: 0 + simulateMessageFromServer(if (softReconnectParam == 0) TestData.JOIN else TestData.RECONNECT) + } + + private fun report(vararg stats: RTCStats) = RTCStatsReport(0, stats.associateBy { it.id }) + + private fun List.single(name: String) = + filter { it.name == name }.also { assertEquals("one $name: ${map { span -> span.name }}", 1, it.size) }.first() + + companion object { + const val COLLECTOR_OUTPUT = "/tmp/livekit-telemetry-otlp.jsonl" + private const val REMOTE_TRACK_ID = "remote_audio" + private const val SEVERITY_WARN = 13 + private val DEVICE_EVENTS = listOf("lk.device.capture.failed", "lk.device.audio_route.changed", "lk.device.audio.interruption") + } +} + +/** What the collector wrote: OTLP/JSON, one export request per line, from [since] (Unix nanoseconds) on. */ +class OtlpFile(file: File, since: Long) { + class Log(val eventName: String, val body: String?, val traceId: String, val spanId: String, val severity: Int, val attributes: Map) + + /** [events] are the span's checkpoints (`ws_open`, `first_media`, `attempt 1 quick`, ...). */ + class Span(val name: String, val traceId: String, val attributes: Map, val events: List) + + val logs = mutableListOf() + val spans = mutableListOf() + + init { + for (line in file.readLines()) { + val request = runCatching { Json.parseToJsonElement(line).jsonObject }.getOrNull() ?: continue + for (scope in request.children("resourceLogs", "scopeLogs")) { + for (r in scope.array("logRecords").map { it.jsonObject }) { + if (r.long("timeUnixNano") < since) continue + logs += Log( + eventName = r.string("eventName") ?: "", + body = r["body"]?.jsonObject?.string("stringValue"), + traceId = r.string("traceId") ?: "", + spanId = r.string("spanId") ?: "", + severity = r["severityNumber"]?.jsonPrimitive?.intOrNull ?: 0, + attributes = attributes(r["attributes"]), + ) + } + } + for (scope in request.children("resourceSpans", "scopeSpans")) { + for (s in scope.array("spans").map { it.jsonObject }) { + if (s.long("startTimeUnixNano") < since) continue + spans += Span( + name = s.string("name") ?: "", + traceId = s.string("traceId") ?: "", + attributes = attributes(s["attributes"]), + events = s.array("events").mapNotNull { it.jsonObject.string("name") }, + ) + } + } + } + } + + private fun JsonObject.children(resources: String, scopes: String) = array(resources).flatMap { it.jsonObject.array(scopes) }.map { it.jsonObject } + + private fun JsonObject.array(key: String): List = this[key]?.jsonArray ?: emptyList() + + private fun JsonObject.string(key: String): String? = (this[key] as? JsonPrimitive)?.contentOrNull + + private fun JsonObject.long(key: String): Long = string(key)?.toLongOrNull() ?: 0L + + /** OTLP/JSON attributes (`[{key, value: {stringValue | intValue | boolValue | doubleValue}}]`) as strings. */ + private fun attributes(value: JsonElement?): Map = + (value?.jsonArray ?: emptyList()).mapNotNull { pair -> + val entry = pair.jsonObject + val key = entry.string("key") ?: return@mapNotNull null + val any = entry["value"]?.jsonObject ?: return@mapNotNull null + val text = any.string("stringValue") + ?: any.string("intValue") + ?: any["boolValue"]?.jsonPrimitive?.booleanOrNull?.toString() + ?: any.string("doubleValue") + ?: return@mapNotNull null + key to text + }.toMap() +} diff --git a/livekit-android-test/src/test/java/io/livekit/android/telemetry/TelemetryPlatformTest.kt b/livekit-android-test/src/test/java/io/livekit/android/telemetry/TelemetryPlatformTest.kt new file mode 100644 index 00000000..ed853f98 --- /dev/null +++ b/livekit-android-test/src/test/java/io/livekit/android/telemetry/TelemetryPlatformTest.kt @@ -0,0 +1,479 @@ +/* + * Copyright 2026 LiveKit, Inc. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package io.livekit.android.telemetry + +import android.media.AudioManager +import com.twilio.audioswitch.AudioDeviceChangeListener +import io.livekit.android.audio.AudioSwitchHandler +import io.livekit.android.events.RoomEvent +import io.livekit.android.room.DefaultsManager +import io.livekit.android.room.network.ReconnectContext +import io.livekit.android.room.network.ReconnectPolicy +import io.livekit.android.room.track.LocalVideoTrack +import io.livekit.android.room.track.LocalVideoTrackOptions +import io.livekit.android.room.track.RemoteTrackPublication +import io.livekit.android.room.track.VideoCaptureParameter +import io.livekit.android.test.MockE2ETest +import io.livekit.android.test.events.EventCollector +import io.livekit.android.test.mock.MockEglBase +import io.livekit.android.test.mock.MockPeerConnection +import io.livekit.android.test.mock.MockRTCThreadToken +import io.livekit.android.test.mock.MockVideoCapturer +import io.livekit.android.test.mock.MockVideoStreamTrack +import io.livekit.android.test.mock.TestData +import io.livekit.android.test.mock.room.track.createMockLocalAudioTrack +import io.livekit.android.util.LKLog +import io.livekit.android.util.LoggingLevel +import io.livekit.uniffi.TelemetryScope +import io.livekit.uniffi.TelemetrySpan +import io.livekit.uniffi.telemetryDisconnectReason +import kotlinx.coroutines.CompletableDeferred +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.ExperimentalCoroutinesApi +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.asContextElement +import kotlinx.coroutines.cancel +import kotlinx.coroutines.cancelAndJoin +import kotlinx.coroutines.channels.Channel +import kotlinx.coroutines.delay +import kotlinx.coroutines.job +import kotlinx.coroutines.launch +import kotlinx.coroutines.runBlocking +import kotlinx.coroutines.test.StandardTestDispatcher +import kotlinx.coroutines.test.advanceUntilIdle +import kotlinx.coroutines.test.runCurrent +import kotlinx.coroutines.withContext +import kotlinx.coroutines.withTimeout +import livekit.LivekitModels +import livekit.LivekitRtc +import livekit.org.webrtc.PeerConnection +import livekit.org.webrtc.RTCStats +import livekit.org.webrtc.RTCStatsCollectorCallback +import livekit.org.webrtc.RTCStatsReport +import okhttp3.mockwebserver.MockResponse +import okhttp3.mockwebserver.MockWebServer +import org.junit.Assert.assertArrayEquals +import org.junit.Assert.assertEquals +import org.junit.Assert.assertFalse +import org.junit.Assert.assertNull +import org.junit.Assert.assertTrue +import org.junit.Assume.assumeTrue +import org.junit.Test +import org.junit.runner.RunWith +import org.mockito.kotlin.any +import org.mockito.kotlin.anyOrNull +import org.mockito.kotlin.argumentCaptor +import org.mockito.kotlin.atLeastOnce +import org.mockito.kotlin.clearInvocations +import org.mockito.kotlin.doAnswer +import org.mockito.kotlin.doReturn +import org.mockito.kotlin.inOrder +import org.mockito.kotlin.mock +import org.mockito.kotlin.never +import org.mockito.kotlin.times +import org.mockito.kotlin.verify +import org.mockito.kotlin.verifyNoInteractions +import org.mockito.kotlin.whenever +import org.mockito.stubbing.Answer +import org.robolectric.RobolectricTestRunner +import uniffi.livekit_telemetry.DeviceEvent +import uniffi.livekit_telemetry.ExportException +import uniffi.livekit_telemetry.ExportRequest +import uniffi.livekit_telemetry.InternalException +import uniffi.livekit_telemetry.SpanName +import java.util.concurrent.CountDownLatch +import java.util.concurrent.TimeUnit +import kotlin.concurrent.thread +import kotlin.time.Duration +import uniffi.livekit_telemetry.DisconnectReason as FfiDisconnectReason + +/** The platform paths the end-to-end test cannot see: tokens, redirects, poll pacing, the opt-out's cache, shared audio listeners. */ +@ExperimentalCoroutinesApi +@RunWith(RobolectricTestRunner::class) +class TelemetryPlatformTest : MockE2ETest() { + + @Test + fun everyTokenReachesTheScope() = runTest { + val scope = mock() + component.rtcEngine().telemetryScope = scope + + connect() + simulateMessageFromServer(TestData.REFRESH_TOKEN) + + inOrder(scope) { + verify(scope).setServer(TestData.EXAMPLE_URL, "token") + verify(scope).setServer(TestData.EXAMPLE_URL, TestData.REFRESH_TOKEN.refreshToken) + } + } + + @Test + fun wakesFasterThanTheIntervalNeverStarveThePolls() = runTest { + val wake = Channel(Channel.CONFLATED) + var polls = 0 + val poller = launch { pollStats(interval = { 1000 }, wake = wake, now = { testScheduler.currentTime }) { polls++ } } + repeat(50) { // a track event every 100 ms, for 5 s + wake.trySend(Unit) + delay(100) + } + poller.cancel() + assertTrue("one poll per second whatever wakes it: $polls", polls >= 4) + } + + @Test + fun aConnectedRoomCollectsNothingAfterTheOptOut() = runTest { + val wasDisabled = Telemetry.disabled + Telemetry.disabled = false // an earlier test in this JVM may have opted out + connect() + val scope = mock() + doAnswer { 50L }.whenever(scope).statsPollIntervalMs() + val subscriber = getSubscriberPeerConnection() + subscriber.statsReport = RTCStatsReport(0, mapOf("IN" to RTCStats(0, "inbound-rtp", "IN", mapOf("bytesReceived" to 1L)))) + val connection = CoroutineScope(SupervisorJob()) + RTCTelemetry(room, scope).start(connection) + try { + repeat(100) { if (subscriber.statsRequests == 0) Thread.sleep(50) } + assertTrue("the Room's stats are collected while telemetry is on", subscriber.statsRequests > 0) + + Telemetry.disabled = true // what disableTelemetry() sets before it returns + Thread.sleep(200) // a poll already running finishes + val collected = subscriber.statsRequests + clearInvocations(scope) + Thread.sleep(500) // ten poll intervals + simulateMessageFromServer(TestData.PARTICIPANT_JOIN) // remote tracks after the opt-out + + assertEquals("no getStats() after the opt-out", collected, subscriber.statsRequests) + verifyNoInteractions(scope) + } finally { + Telemetry.disabled = wasDisabled + connection.cancel() + } + } + + @Test + fun anOptOutDuringAPollStartsNoFurtherRequestAndSubmitsNothing() = runTest { + val wasDisabled = Telemetry.disabled + Telemetry.disabled = false // an earlier test in this JVM may have opted out + // fastPublish: the publisher exists from the join, so a poll asks it first, then the subscriber. + connect(TestData.JOIN.toBuilder().apply { join = join.toBuilder().setFastPublish(true).build() }.build()) + val scope = mock() + val publisher = getPublisherPeerConnection() + val subscriber = getSubscriberPeerConnection() + val report = RTCStatsReport(0, mapOf("IN" to RTCStats(0, "inbound-rtp", "IN", mapOf("bytesReceived" to 1L)))) + val rtc = RTCTelemetry(room, scope) // not started: the test runs one poll itself + + /** One poll, paused in [pc]'s answer while the opt-out happens, then run to its end. */ + fun pollOptingOutDuring(pc: MockPeerConnection) = runBlocking(Dispatchers.Default) { + val held = CompletableDeferred() + pc.statsAnswer = { held.complete(it) } + val poll = launch { rtc.recordPeerStats() } + val answer = withTimeout(5_000) { held.await() } + synchronized(Telemetry) { Telemetry.disabled = true } // disableTelemetry()'s platform half; the core's opt-out would last the JVM + answer.onStatsDelivered(report) + poll.join() + } + + try { + // Paused before the subscriber's request is initiated: it never is. + pollOptingOutDuring(publisher) + assertEquals("no getStats() begins after the opt-out", 0, subscriber.statsRequests) + + // Paused before the submit (both answers in hand): nothing reaches the core. + Telemetry.disabled = false + publisher.statsAnswer = { it.onStatsDelivered(report) } + pollOptingOutDuring(subscriber) + verify(scope, never()).recordPeerStats(any(), any(), anyOrNull()) + } finally { + Telemetry.disabled = wasDisabled + } + } + + @Test + fun aThrowingCoreNeverFailsLoggingDisconnectOrPublishCleanup() = runTest { + val panic = Answer { throw InternalException("core panic") } + val span = mock(defaultAnswer = panic) + val scope = mock(defaultAnswer = panic) + doReturn(span).whenever(scope).start(any(), anyOrNull()) + val installed = Telemetry.installed + val wasDisabled = Telemetry.disabled + Telemetry.installed = true // so logs and device events reach the (throwing) core, + Telemetry.disabled = false // whatever an earlier test in this JVM did + // The app's logger throws on the very diagnostic that reports the core's failure. + val appLogger = LKLog.logger + val appLevel = LKLog.loggingLevel + LKLog.loggingLevel = LoggingLevel.WARN + var loggerThrew = 0 + LKLog.logger = object : LKLog.Logger { + override fun log(priority: LoggingLevel, t: Throwable?, message: String) { + if (t is InternalException) { + loggerThrew++ + throw IllegalStateException("app logger failed") + } + appLogger?.log(priority, t, message) + } + } + try { + room.telemetryScope = scope + component.rtcEngine().telemetryScope = scope + room.localParticipant.telemetryScope = scope + connect(TestData.JOIN.toBuilder().apply { join = join.toBuilder().setFastPublish(true).build() }.build()) + + // An SDK warning from a Room handler (outside the telemetry package, which capture skips), + // filed under the Room's throwing scope; and a device event. + simulateMessageFromServer( + LivekitRtc.SignalResponse.newBuilder() + .setTrackUnpublished(LivekitRtc.TrackUnpublishedResponse.newBuilder().setTrackSid("TR_unknown")) + .build(), + ) + verify(scope, atLeastOnce()).log(any()) + Telemetry.deviceEvent(DeviceEvent.AudioInterruption(began = true)) + + // A publish the server never answers, cancelled: the video transceiver is still rolled back. + wsFactory.unregisterSignalRequestHandler(wsFactory.defaultSignalRequestHandler) + val publish = launch(StandardTestDispatcher(testScheduler)) { room.localParticipant.publishVideoTrack(createVideoTrack()) } + runCurrent() + publish.cancelAndJoin() + val transceiver = getPublisherPeerConnection().transceivers.single() + verify(transceiver).stopInternal() + + // Disconnect: the engine still closes, and the Room says it disconnected. + val subscriber = getSubscriberPeerConnection() + val events = EventCollector(room.events, coroutineRule.scope) + room.disconnect() + assertEquals(PeerConnection.IceConnectionState.CLOSED, subscriber.iceConnectionState()) + assertTrue(events.stopCollecting().any { it is RoomEvent.Disconnected }) + assertTrue("the app's logger threw on the failure diagnostic", loggerThrew > 0) + } finally { + Telemetry.installed = installed + Telemetry.disabled = wasDisabled + LKLog.logger = appLogger + LKLog.loggingLevel = appLevel + } + } + + private fun createVideoTrack() = LocalVideoTrack( + capturer = MockVideoCapturer(), + source = mock(), + name = "", + options = LocalVideoTrackOptions(isScreencast = false, deviceId = null, position = null, captureParams = VideoCaptureParameter(1280, 720, 30)), + rtcTrack = MockVideoStreamTrack(), + peerConnectionFactory = component.peerConnectionFactory(), + context = context, + eglBase = MockEglBase(), + defaultsManager = DefaultsManager(), + trackFactory = mock(), + rtcThreadToken = MockRTCThreadToken(), + ) + + @Test + fun optOutDeletesAPreviousLaunchsCache() { + val cache = Telemetry.storageDirectory(context).apply { mkdirs() } + cache.resolve("batch").writeText("unsent") + Telemetry.disabled = true + try { + assertNull("no scope after the opt-out", Telemetry.scope(context)) + assertFalse("the cache left by a previous launch is gone", cache.exists()) + } finally { + Telemetry.disabled = false + } + } + + @Test + fun aSharedAudioHandlerHasOneListenerPairUntilItsLastRoomLeaves() { + // An app-supplied handler is shared by every Room; route and focus are process events. + val handler = mock() + val first = handler.observeForTelemetry() + val second = handler.observeForTelemetry() + val route = argumentCaptor() + val focus = argumentCaptor() + verify(handler, times(1)).registerAudioDeviceChangeListener(route.capture()) + verify(handler, times(1)).registerOnAudioFocusChangeListener(focus.capture()) + + first() + first() // a release is idempotent + verify(handler, never()).unregisterAudioDeviceChangeListener(any()) + second() + verify(handler).unregisterAudioDeviceChangeListener(route.firstValue) + verify(handler).unregisterOnAudioFocusChangeListener(focus.firstValue) + + val third = handler.observeForTelemetry() // a later Room registers afresh + verify(handler, times(2)).registerAudioDeviceChangeListener(any()) + third() + verify(handler, times(2)).unregisterAudioDeviceChangeListener(any()) + } + + @Test + fun aRoomCreatedFromAnAudioCallbackNeverDeadlocksTheLastRelease() { + // The handler dispatches holding its listener lock; a callback there may create a Room + // while another thread releases the last Room, whose unregister needs that lock. + val listenerLock = Object() + val unregistering = CountDownLatch(1) + val handler = mock() + doAnswer { synchronized(listenerLock) {} }.whenever(handler).registerAudioDeviceChangeListener(any()) + doAnswer { + unregistering.countDown() + synchronized(listenerLock) {} + }.whenever(handler).unregisterAudioDeviceChangeListener(any()) + val last = handler.observeForTelemetry() + var created: (() -> Unit)? = null + + val dispatch = thread { + synchronized(listenerLock) { + unregistering.await(5, TimeUnit.SECONDS) + created = handler.observeForTelemetry() // a Room created from the route callback + } + } + val release = thread { last() } + dispatch.join(5_000) + release.join(5_000) + + assertFalse("no deadlock", dispatch.isAlive || release.isAlive) + // The release unregistered, then re-registered for the Room the callback created. + verify(handler, times(2)).registerAudioDeviceChangeListener(any()) + created!!() + verify(handler, times(2)).unregisterAudioDeviceChangeListener(any()) + } + + @Test + fun aPublishCancelledBeforeItStartsLeavesNoOpenSpan() = runTest { + val span = mock() + val scope = mock() + doReturn(span).whenever(scope).start(any(), anyOrNull()) + room.localParticipant.telemetryScope = scope + launch { + coroutineContext.job.cancel() // wins before withContext runs the publish + room.localParticipant.publishAudioTrack(createMockLocalAudioTrack()) + }.join() + verify(scope).start(SpanName.Publish, null) + verify(span).cancel() + } + + @Test + fun aConnectAttemptHandsOverItsServerBeforeTheJoin() = runTest { + val scope = mock() + room.telemetryScope = scope // the engine has none: only Room.connect can call it + connect() + verify(scope).setServer(TestData.EXAMPLE_URL, "token") + } + + @Test + fun aPublishFilesItsWarningsUnderItsOwnSpanAndSkipsAnEndedParent() = runTest { + val ended = mock() + doReturn(true).whenever(ended).isEnded() + val span = mock() + val scope = mock() + doReturn(span).whenever(scope).start(any(), anyOrNull()) + val installed = Telemetry.installed + val wasDisabled = Telemetry.disabled + Telemetry.installed = true + Telemetry.disabled = false + room.localParticipant.telemetryScope = scope + try { + // Not connected, so the publish fails its permission check with a warning. + withContext(Telemetry.currentSpan.asContextElement(ended)) { + runCatching { room.localParticipant.publishAudioTrack(createMockLocalAudioTrack()) } + } + verify(scope).start(SpanName.Publish, null) + verify(span, atLeastOnce()).context() // the warning asked lk.publish for its span id + } finally { + Telemetry.installed = installed + Telemetry.disabled = wasDisabled + } + } + + @Test + fun aReconnectThatGivesUpEndsTheSessionAsReconnectFailed() = runTest { + val scope = mock() + room.telemetryScope = scope + connect() + component.rtcEngine().reconnectPolicy = object : ReconnectPolicy { + override fun getNextRetryDelay(context: ReconnectContext): Duration? = null + } + disconnectPeerConnection() + testScheduler.advanceUntilIdle() + verify(scope).disconnected(FfiDisconnectReason.RECONNECT_FAILED) + } + + @Test + fun aServerLeaveKeepsItsProtocolReason() = runTest { + assumeTrue("needs the core's protocol mapping", runCatching { telemetryDisconnectReason(0) }.isSuccess) + val scope = mock() + room.telemetryScope = scope + connect() + simulateMessageFromServer( + LivekitRtc.SignalResponse.newBuilder() + .setLeave(LivekitRtc.LeaveRequest.newBuilder().setReason(LivekitModels.DisconnectReason.AGENT_ERROR)) + .build(), + ) + verify(scope).disconnected(FfiDisconnectReason.AGENT_ERROR) // the SDK's own enum has no agent_error + } + + @Test + fun aManualSubscribeWakesThePollerAtOnce() = runTest { + connect() + simulateMessageFromServer(TestData.PARTICIPANT_JOIN) + val scope = mock() + var interval = 60_000L + doAnswer { interval }.whenever(scope).statsPollIntervalMs() + val subscriber = getSubscriberPeerConnection() + subscriber.statsReport = RTCStatsReport(0, mapOf("IN" to RTCStats(0, "inbound-rtp", "IN", mapOf("bytesReceived" to 1L)))) + val rtc = RTCTelemetry(room, scope) + component.rtcEngine().client.rtcTelemetry = rtc + val wasDisabled = Telemetry.disabled + Telemetry.disabled = false + val connection = CoroutineScope(SupervisorJob()) + rtc.start(connection) + try { + Thread.sleep(300) + assertEquals("idle: the next poll is a minute away", 0, subscriber.statsRequests) + + interval = 50 // what the core answers once a subscribe awaits first media + val publication = room.remoteParticipants.values.single().trackPublications.values.first() as RemoteTrackPublication + publication.setSubscribed(true) + repeat(40) { if (subscriber.statsRequests == 0) Thread.sleep(50) } + + verify(scope).subscribeStarted(any()) + assertTrue("polled right after the manual subscribe", subscriber.statsRequests > 0) + } finally { + Telemetry.disabled = wasDisabled + connection.cancel() + } + } + + @Test + fun transportPassesAnswersThroughAndFollowsNoRedirect() = runTest { + val server = MockWebServer() + server.start() + val transport = OkHttpTelemetryTransport() + val request = ExportRequest(server.url("/v1/logs").toString(), mapOf("Authorization" to "Bearer token"), byteArrayOf(1, 2, 3)) + + server.enqueue(MockResponse().setResponseCode(429).setHeader("Retry-After", "7").setBody("quota")) + val throttled = transport.send(request) + assertEquals(429, throttled.status.toInt()) + assertEquals("7", throttled.headers["Retry-After"]) + assertArrayEquals("quota".toByteArray(), throttled.body) + assertEquals("Bearer token", server.takeRequest().getHeader("Authorization")) + + server.enqueue(MockResponse().setResponseCode(307).setHeader("Location", server.url("/elsewhere"))) + assertEquals("a redirect is the answer, never followed", 307, transport.send(request).status.toInt()) + assertEquals("the redirect target is never requested", 2, server.requestCount) + + server.shutdown() + val noAnswer = runCatching { transport.send(request) }.exceptionOrNull() + assertTrue("no answer is the transport's only error: $noAnswer", noAnswer is ExportException.Retryable) + } +} diff --git a/livekit-android-test/src/test/resources/telemetry/otelcol-lgtm.yaml b/livekit-android-test/src/test/resources/telemetry/otelcol-lgtm.yaml new file mode 100644 index 00000000..6b74fac8 --- /dev/null +++ b/livekit-android-test/src/test/resources/telemetry/otelcol-lgtm.yaml @@ -0,0 +1,11 @@ +# Layered on otelcol.yaml for local testing: also forwards every batch to a local LGTM backend. +# otelcol-contrib --config otelcol.yaml --config otelcol-lgtm.yaml +exporters: + otlphttp: + endpoint: http://localhost:4318 +service: + pipelines: + logs: + exporters: [file, otlphttp] + traces: + exporters: [file, otlphttp] diff --git a/livekit-android-test/src/test/resources/telemetry/otelcol.yaml b/livekit-android-test/src/test/resources/telemetry/otelcol.yaml new file mode 100644 index 00000000..a76b3a11 --- /dev/null +++ b/livekit-android-test/src/test/resources/telemetry/otelcol.yaml @@ -0,0 +1,22 @@ +# Collector for the telemetry end-to-end test (TelemetryMockE2ETest): accepts OTLP/HTTP on :4319 and +# writes every request as one JSON line the test reads back. Run: +# otelcol-contrib --config livekit-android-test/src/test/resources/telemetry/otelcol.yaml +receivers: + otlp: + protocols: + http: + endpoint: 127.0.0.1:4319 +exporters: + file: + path: /tmp/livekit-telemetry-otlp.jsonl +service: + telemetry: + metrics: + level: none + pipelines: + logs: + receivers: [otlp] + exporters: [file] + traces: + receivers: [otlp] + exporters: [file] diff --git a/settings.gradle b/settings.gradle index b5882cbb..909579af 100644 --- a/settings.gradle +++ b/settings.gradle @@ -11,9 +11,27 @@ pluginManagement { gradlePluginPortal() } } +// Development wiring for unreleased Rust bindings: `-PlivekitUniffiVersion=` (or the +// LIVEKIT_UNIFFI_VERSION environment variable) resolves livekit-uniffi-android at that version from +// Maven Local, where `cargo make android-package-local` in rust-sdks/livekit-uniffi publishes it. +// Unset, the released artifact from gradle/libs.versions.toml is used. +def localUniffi = providers.gradleProperty('livekitUniffiVersion') + .orElse(providers.environmentVariable('LIVEKIT_UNIFFI_VERSION')).getOrNull() +if (localUniffi) { + gradle.beforeProject { project -> + project.configurations.configureEach { + resolutionStrategy.eachDependency { + if (requested.group == 'io.livekit' && requested.name == 'livekit-uniffi-android') useVersion(localUniffi) + } + } + } +} dependencyResolutionManagement { repositoriesMode.set(RepositoriesMode.FAIL_ON_PROJECT_REPOS) repositories { + if (localUniffi) { + mavenLocal { content { includeModule('io.livekit', 'livekit-uniffi-android') } } + } google() mavenCentral() maven { url 'https://jitpack.io' }