package com.suno.android.media import androidx.media3.common.MediaItem import androidx.media3.common.Player import com.suno.android.common_analytics.listening_source.ListeningSourceCache import com.suno.android.common_analytics.listening_source.toListeningSourceContext import com.suno.android.common_analytics.managers.AnalyticsManager import com.suno.android.common_analytics.managers.Interval import com.suno.android.common_analytics.managers.NewSongCause import com.suno.android.common_core_utils.ApplicationCoroutineScope import com.suno.android.common_core_utils.Id import com.suno.android.common_core_utils.SunoLogger import com.suno.android.common_data.repos.ClipMetricsRepository import com.suno.android.common_data.user.UserSessionRepository import com.suno.android.common_networking.remote.session.User import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.channels.BufferOverflow import kotlinx.coroutines.flow.MutableSharedFlow import kotlinx.coroutines.flow.catch import kotlinx.coroutines.flow.debounce import kotlinx.coroutines.flow.flowOn import kotlinx.coroutines.flow.launchIn import kotlinx.coroutines.flow.onEach import kotlinx.coroutines.launch import java.util.concurrent.ConcurrentHashMap import javax.inject.Inject import javax.inject.Singleton import kotlin.time.Duration import kotlin.time.Duration.Companion.seconds interface MediaAnalyticsManager { fun registerPlayer( player: Player, ) fun unregisterPlayer( player: Player, ) } internal class WatchdogTimer( logger: SunoLogger, coroutineScope: CoroutineScope, timeout: Duration, private val onTimeout: suspend () -> Unit, ) { private val petSignal = MutableSharedFlow( replay = 0, extraBufferCapacity = 1, onBufferOverflow = BufferOverflow.DROP_OLDEST, ) private val job = petSignal.debounce(timeout) .onEach { onTimeout() } .catch { exception -> logger.e(exception) } .launchIn(coroutineScope) fun pet() { petSignal.tryEmit(Unit) } fun cancel() { job.cancel() } } @Singleton internal class MediaAnalyticsManagerImpl @Inject constructor( @ApplicationCoroutineScope private val coroutineScope: CoroutineScope, private val analyticsManager: AnalyticsManager, private val listeningSourceCache: ListeningSourceCache, private val clipMetricsRepository: ClipMetricsRepository, private val userSessionRepository: UserSessionRepository, loggerFactory: SunoLogger.Factory, ) : MediaAnalyticsManager { private val logger = loggerFactory.create(this@MediaAnalyticsManagerImpl) private val registeredPlayers = ConcurrentHashMap() private val scope = CoroutineScope(Dispatchers.Default) @Volatile private var currentUser: User? = null init { updateCurrentUser() } private fun updateCurrentUser() { userSessionRepository.sessionConfigurationStateFlow() .onEach { sessionConfiguration -> currentUser = sessionConfiguration.user } .catch { exception -> logger.e(exception) } .flowOn(Dispatchers.IO) .launchIn(coroutineScope) } override fun registerPlayer( player: Player, ) { if (registeredPlayers.containsKey(player)) { return } logger.d { "Registering player: ${player.hashCode()}" } val watchdogTimer = WatchdogTimer( logger = logger, coroutineScope = coroutineScope, timeout = 2.seconds, onTimeout = { flushScrub(player, sendEvents = true) }, ) val listener = createPlayerListener(player).also(player::addListener) registeredPlayers[player] = RegisteredPlayerState( listener = listener, watchdogTimer = watchdogTimer, analyticsState = MediaAnalyticsState(), ) } private fun flushScrub( player: Player, sendEvents: Boolean, ) { val analyticsState = getAnalyticsState(player) ?: return if (!analyticsState.isScrubbing) { return } if (sendEvents && analyticsState.isPlaying) { val interval = Interval(analyticsState.continuousPlayStart, analyticsState.scrubBounds.startTime) analyticsManager.trackSeekProgressBarPauseSong(interval) analyticsManager.trackSeekProgressBarPlaySong() } updatePlayerState(player) { it.copy( isScrubbing = false, continuousPlayStart = it.scrubBounds.endTime, ) } } private fun createPlayerListener( attachedToPlayer: Player, ): Player.Listener { return object : Player.Listener { private val logPrefix = "Player(${attachedToPlayer.hashCode()})" override fun onPlayWhenReadyChanged( playWhenReady: Boolean, reason: Int, ) { super.onPlayWhenReadyChanged(playWhenReady, reason) logger.d { "$logPrefix - playWhenReadyChanged: ${playWhenReadyChangeReasonToString(reason)}" } updatePlayerState(attachedToPlayer) { state -> state.copy( isPlaying = playWhenReady, playWhenReadyChangedReason = fromReasonCode(reason) ?: PlayWhenReadyChangedReason.USER_REQUEST, ) } } override fun onMediaItemTransition( mediaItem: MediaItem?, reason: Int, ) { super.onMediaItemTransition(mediaItem, reason) logger.d { "$logPrefix - onMediaItemTransition: ${mediaItemTransitionReasonToString(reason)}" } val analyticsState = getAnalyticsState(attachedToPlayer) ?: return val reasonAsEnum = fromReasonCode(reason) if (analyticsState.ignoreRedundantSeek && reasonAsEnum == MediaItemTransitionReason.SEEK) { return } updatePlayerState(attachedToPlayer) { state -> state.copy( mediaItemTransitionReason = reasonAsEnum ?: MediaItemTransitionReason.PLAYLIST_CHANGED, lastPlayingClipIndex = state.nowPlayingClipIndex, nowPlayingClipIndex = attachedToPlayer.currentMediaItemIndex, ignoreRedundantSeek = reasonAsEnum == MediaItemTransitionReason.PLAYLIST_CHANGED, ) } } override fun onPositionDiscontinuity( oldPosition: Player.PositionInfo, newPosition: Player.PositionInfo, reason: Int, ) { super.onPositionDiscontinuity(oldPosition, newPosition, reason) logger.d { "$logPrefix - onPositionDiscontinuity: ${positionDiscontinuityReasonToString(reason)}" } val reasonAsEnum = fromReasonCode(reason) val currentMediaItem = attachedToPlayer.currentMediaItem if (reasonAsEnum == null || currentMediaItem == null) return updatePlayerState(attachedToPlayer) { state -> val discontinuity = reasonAsEnum to Interval(oldPosition.positionMs, newPosition.positionMs) state.copy( playbackDiscontinuities = state.playbackDiscontinuities + discontinuity, playbackDiscontinuityReason = reasonAsEnum, ) } val isSeekRelated = reasonAsEnum == PlaybackDiscontinuityReason.SEEK || reasonAsEnum == PlaybackDiscontinuityReason.SEEK_ADJUSTMENT if (isSeekRelated) { handleSeekDiscontinuity(attachedToPlayer, oldPosition, newPosition, currentMediaItem) } } override fun onEvents( player: Player, events: Player.Events, ) { super.onEvents(player, events) logger.d { "$logPrefix - onEvents: ${events.toDebugString()}" } if (events.contains(Player.EVENT_MEDIA_ITEM_TRANSITION)) { attachedToPlayer.currentMediaItem?.let { mediaItem -> handleMediaItemTransition(attachedToPlayer, mediaItem) } } val analyticsState = getAnalyticsState(attachedToPlayer) ?: return val hasPlayWhenReadyChanged = events.contains(Player.EVENT_PLAY_WHEN_READY_CHANGED) when { analyticsState.isPlaying && analyticsState.deferredPlayEventCause != null -> { handleMediaItemTransitionCompleted( attachedToPlayer, hasPlayWhenReadyChanged, analyticsState.deferredPlayEventCause, ) } hasPlayWhenReadyChanged -> handlePlayPause(attachedToPlayer) events.contains(Player.EVENT_PLAYBACK_STATE_CHANGED) && player.playbackState == Player.STATE_ENDED -> { handlePlaylistEnded(attachedToPlayer) } } } } } private fun handleSeekDiscontinuity( player: Player, oldPosition: Player.PositionInfo, newPosition: Player.PositionInfo, currentMediaItem: MediaItem, ) { val playerState = registeredPlayers[player] ?: return val analyticsState = playerState.analyticsState when { analyticsState.isScrubbing -> handleScrubbing(player, newPosition.positionMs, playerState) newPosition.positionMs == 0L && oldPosition.positionMs > 0L -> handleSeekToStart( player, oldPosition.positionMs, currentMediaItem, ) else -> startScrubbing(player, oldPosition.positionMs, newPosition.positionMs, playerState) } } private fun handleMediaItemTransition( player: Player, newMediaItem: MediaItem, ) { updatePlayerState(player) { it.copy(ignoreRedundantSeek = false) } flushScrub(player, sendEvents = false) val analyticsState = getAnalyticsState(player) ?: return val (cause, discontinuityReason) = analyticsState.toTransitionCause() val discontinuityInterval = analyticsState.playbackDiscontinuities.getOrDefault( discontinuityReason, Interval(0, 0), ) val deferredPlayEventUpdate = createDeferredPlayEventUpdate( newMediaItem, cause, Interval(analyticsState.continuousPlayStart, discontinuityInterval.startTime), discontinuityInterval.endTime, ) updatePlayerState(player, deferredPlayEventUpdate) } private fun handleMediaItemTransitionCompleted( player: Player, hasPlayWhenReadyChanged: Boolean, cause: NewSongCause, ) { val analyticsState = getAnalyticsState(player) ?: return val listeningSource = listeningSourceCache .getAndRemove(Id(analyticsState.deferredPlayEventClipId)) ?.toListeningSourceContext() if (!hasPlayWhenReadyChanged && !analyticsState.wasAtEnd) { analyticsManager.trackPlayNewSongPauseSong( interval = analyticsState.deferredPlayEventPauseInterval, cause = cause, listeningSource = listeningSource, ) } analyticsManager.trackPlayNewSong( clipId = analyticsState.deferredPlayEventClipId, cause = cause, isOwned = analyticsState.deferredPlayEventClipIsOwned, listeningSource = listeningSource, ) scope.launch { clipMetricsRepository.incrementPlayCount(clipId = analyticsState.deferredPlayEventClipId) } updatePlayerState(player) { it.copy( deferredPlayEventCause = null, wasAtEnd = false, ) } } private fun handlePlayPause( player: Player, ) { val analyticsState = getAnalyticsState(player) ?: return if (analyticsState.isPlaying) { analyticsManager.trackPlaySong() } else { val currentPosition = player.currentPosition val interval = Interval(analyticsState.continuousPlayStart, currentPosition) analyticsManager.trackPauseSong(interval) updatePlayerState(player) { it.copy(continuousPlayStart = currentPosition) } } } private fun handlePlaylistEnded( player: Player, ) { val analyticsState = getAnalyticsState(player) ?: return val interval = Interval(analyticsState.continuousPlayStart, player.duration) analyticsManager.trackSongEnd(interval) updatePlayerState(player) { it.copy(wasAtEnd = true) } } private fun handleScrubbing( player: Player, newPositionMs: Long, playerState: RegisteredPlayerState, ) { updatePlayerState(player) { state -> state.copy( scrubBounds = state.scrubBounds.copy( state.scrubBounds.startTime, newPositionMs, ), ) } playerState.watchdogTimer.pet() } private fun handleSeekToStart( player: Player, oldPositionMs: Long, currentMediaItem: MediaItem, ) { val analyticsState = getAnalyticsState(player) ?: return val deferredPlayEventUpdate = createDeferredPlayEventUpdate( currentMediaItem, NewSongCause.SEEK_TO_START, Interval(analyticsState.continuousPlayStart, oldPositionMs), continuousPlayStart = 0L, ) updatePlayerState(player, deferredPlayEventUpdate) } private fun startScrubbing( player: Player, oldPositionMs: Long, newPositionMs: Long, playerState: RegisteredPlayerState, ) { updatePlayerState(player) { it.copy( isScrubbing = true, scrubBounds = Interval(oldPositionMs, newPositionMs), ) } playerState.watchdogTimer.pet() } override fun unregisterPlayer( player: Player, ) { logger.d { "Un-registered player: ${player.hashCode()}" } val playerState = registeredPlayers.remove(player) playerState?.let { state -> player.removeListener(state.listener) state.watchdogTimer.cancel() } } private fun getAnalyticsState( player: Player, ): MediaAnalyticsState? = registeredPlayers[player]?.analyticsState private fun updatePlayerState( player: Player, updateFunction: (MediaAnalyticsState) -> MediaAnalyticsState, ) { registeredPlayers.computeIfPresent(player) { _, playerState -> playerState.copy( analyticsState = updateFunction(playerState.analyticsState), ) } } private fun createDeferredPlayEventUpdate( mediaItem: MediaItem, cause: NewSongCause, pauseInterval: Interval, continuousPlayStart: Long, ): (MediaAnalyticsState) -> MediaAnalyticsState = { state -> state.copy( deferredPlayEventCause = cause, deferredPlayEventPauseInterval = pauseInterval, deferredPlayEventClipId = mediaItem.mediaId, deferredPlayEventClipIsOwned = mediaItem.isOwnedBy(currentUser), continuousPlayStart = continuousPlayStart, ) } private fun MediaAnalyticsState.toTransitionCause(): Pair = when (mediaItemTransitionReason) { MediaItemTransitionReason.SEEK -> { val mediaIndexDelta = nowPlayingClipIndex - lastPlayingClipIndex Pair( when (mediaIndexDelta) { 1 -> NewSongCause.SKIP_FORWARD -1 -> NewSongCause.SKIP_BACKWARD else -> NewSongCause.VANILLA }, PlaybackDiscontinuityReason.SEEK, ) } MediaItemTransitionReason.AUTO -> Pair( NewSongCause.AUTO_ADVANCE, PlaybackDiscontinuityReason.AUTO_TRANSITION, ) MediaItemTransitionReason.PLAYLIST_CHANGED -> Pair( NewSongCause.VANILLA, PlaybackDiscontinuityReason.REMOVE, ) MediaItemTransitionReason.REPEAT -> Pair( NewSongCause.AUTO_REPEAT, PlaybackDiscontinuityReason.AUTO_TRANSITION, ) } private fun MediaItem.isOwnedBy( user: User?, ): Boolean = user?.let { mediaMetadata.extras?.getString("artistUserId") == it.id } == true } private data class RegisteredPlayerState( val listener: Player.Listener, val watchdogTimer: WatchdogTimer, val analyticsState: MediaAnalyticsState, ) data class MediaAnalyticsState( val continuousPlayStart: Long = 0, val isPlaying: Boolean = false, val playbackDiscontinuityReason: PlaybackDiscontinuityReason = PlaybackDiscontinuityReason.AUTO_TRANSITION, val playbackDiscontinuities: Map = emptyMap(), val playWhenReadyChangedReason: PlayWhenReadyChangedReason = PlayWhenReadyChangedReason.USER_REQUEST, val isScrubbing: Boolean = false, val scrubBounds: Interval = Interval(0, 0), val mediaItemTransitionReason: MediaItemTransitionReason = MediaItemTransitionReason.AUTO, val lastPlayingClipIndex: Int = 0, val nowPlayingClipIndex: Int = 0, val deferredPlayEventCause: NewSongCause? = null, val deferredPlayEventPauseInterval: Interval = Interval(0, 0), val deferredPlayEventClipId: String = "", val deferredPlayEventClipIsOwned: Boolean = false, val ignoreRedundantSeek: Boolean = false, val wasAtEnd: Boolean = false, )