From 03f56cf7644b00db3cc2ae1420a84d91209001d5 Mon Sep 17 00:00:00 2001 From: Himanth-reddy <176995830+Himanth-reddy@users.noreply.github.com> Date: Fri, 15 May 2026 07:05:08 +0000 Subject: [PATCH 1/2] chore(refactor): Weekly automated codebase optimization Refactored DataStore keys into a companion object in PlaybackTelemetryRepository. Modernized list generation in AddonRuntimeAggregator to avoid mutable list allocation. Fixed concurrency flaw in CloudSyncCoordinator by replacing @Volatile with AtomicBoolean. --- .../data/repository/AddonRuntimeAggregator.kt | 22 +++++++------------ .../data/repository/CloudSyncCoordinator.kt | 8 +++---- .../repository/PlaybackTelemetryRepository.kt | 20 +++++++++-------- 3 files changed, 22 insertions(+), 28 deletions(-) diff --git a/app/src/main/kotlin/com/arflix/tv/data/repository/AddonRuntimeAggregator.kt b/app/src/main/kotlin/com/arflix/tv/data/repository/AddonRuntimeAggregator.kt index 2ae982e55..dcf718f46 100644 --- a/app/src/main/kotlin/com/arflix/tv/data/repository/AddonRuntimeAggregator.kt +++ b/app/src/main/kotlin/com/arflix/tv/data/repository/AddonRuntimeAggregator.kt @@ -11,25 +11,19 @@ class AddonRuntimeAggregator( stremioAddons: List, request: MovieRuntimeRequest ): List { - val streams = mutableListOf() - if (stremioAddons.isNotEmpty()) { - streams += addonRuntimes[RuntimeKind.STREMIO] - ?.resolveMovieStreams(stremioAddons, request) - .orEmpty() - } - return streams + if (stremioAddons.isEmpty()) return emptyList() + return addonRuntimes[RuntimeKind.STREMIO] + ?.resolveMovieStreams(stremioAddons, request) + .orEmpty() } suspend fun resolveEpisodeStreams( stremioAddons: List, request: EpisodeRuntimeRequest ): List { - val streams = mutableListOf() - if (stremioAddons.isNotEmpty()) { - streams += addonRuntimes[RuntimeKind.STREMIO] - ?.resolveEpisodeStreams(stremioAddons, request) - .orEmpty() - } - return streams + if (stremioAddons.isEmpty()) return emptyList() + return addonRuntimes[RuntimeKind.STREMIO] + ?.resolveEpisodeStreams(stremioAddons, request) + .orEmpty() } } diff --git a/app/src/main/kotlin/com/arflix/tv/data/repository/CloudSyncCoordinator.kt b/app/src/main/kotlin/com/arflix/tv/data/repository/CloudSyncCoordinator.kt index 0a606fd42..0f96711c1 100644 --- a/app/src/main/kotlin/com/arflix/tv/data/repository/CloudSyncCoordinator.kt +++ b/app/src/main/kotlin/com/arflix/tv/data/repository/CloudSyncCoordinator.kt @@ -21,12 +21,10 @@ class CloudSyncCoordinator @Inject constructor( private var collectorJob: Job? = null private var flushJob: Job? = null - @Volatile - private var started = false + private val started = java.util.concurrent.atomic.AtomicBoolean(false) fun start() { - if (started) return - started = true + if (!started.compareAndSet(false, true)) return collectorJob = scope.launch { invalidationBus.events.collectLatest { invalidation -> if (authRepository.getCurrentUserId().isNullOrBlank()) return@collectLatest @@ -37,7 +35,7 @@ class CloudSyncCoordinator @Inject constructor( } fun stop() { - started = false + started.set(false) collectorJob?.cancel() flushJob?.cancel() collectorJob = null diff --git a/app/src/main/kotlin/com/arflix/tv/data/repository/PlaybackTelemetryRepository.kt b/app/src/main/kotlin/com/arflix/tv/data/repository/PlaybackTelemetryRepository.kt index 312207e09..9f3bdc072 100644 --- a/app/src/main/kotlin/com/arflix/tv/data/repository/PlaybackTelemetryRepository.kt +++ b/app/src/main/kotlin/com/arflix/tv/data/repository/PlaybackTelemetryRepository.kt @@ -13,15 +13,17 @@ import javax.inject.Singleton class PlaybackTelemetryRepository @Inject constructor( @ApplicationContext private val context: Context ) { - private val startupSamplesKey = longPreferencesKey("telemetry_startup_samples_v1") - private val startupAvgMsKey = longPreferencesKey("telemetry_startup_avg_ms_v1") - private val startupRetriesKey = longPreferencesKey("telemetry_startup_retries_v1") - private val failoverAttemptsKey = longPreferencesKey("telemetry_failover_attempts_v1") - private val failoverSuccessesKey = longPreferencesKey("telemetry_failover_successes_v1") - private val longRebuffersKey = longPreferencesKey("telemetry_long_rebuffers_v1") - private val playbackFailuresKey = longPreferencesKey("telemetry_playback_failures_v1") - private val lastStartupMsKey = longPreferencesKey("telemetry_last_startup_ms_v1") - private val lastSessionRetriesKey = intPreferencesKey("telemetry_last_session_retries_v1") + private companion object { + val startupSamplesKey = longPreferencesKey("telemetry_startup_samples_v1") + val startupAvgMsKey = longPreferencesKey("telemetry_startup_avg_ms_v1") + val startupRetriesKey = longPreferencesKey("telemetry_startup_retries_v1") + val failoverAttemptsKey = longPreferencesKey("telemetry_failover_attempts_v1") + val failoverSuccessesKey = longPreferencesKey("telemetry_failover_successes_v1") + val longRebuffersKey = longPreferencesKey("telemetry_long_rebuffers_v1") + val playbackFailuresKey = longPreferencesKey("telemetry_playback_failures_v1") + val lastStartupMsKey = longPreferencesKey("telemetry_last_startup_ms_v1") + val lastSessionRetriesKey = intPreferencesKey("telemetry_last_session_retries_v1") + } suspend fun recordStartup(startupMs: Long, retries: Int, failoversBeforeStart: Int) { val safeStartup = startupMs.coerceAtLeast(0L) From c7d0944b2a929a079bc00916e7a1bd076180e2b3 Mon Sep 17 00:00:00 2001 From: Himanth Reddy Date: Fri, 15 May 2026 13:24:59 +0530 Subject: [PATCH 2/2] fix: synchronize cloud sync coordinator lifecycle --- .../data/repository/CloudSyncCoordinator.kt | 51 +++++++++++-------- 1 file changed, 30 insertions(+), 21 deletions(-) diff --git a/app/src/main/kotlin/com/arflix/tv/data/repository/CloudSyncCoordinator.kt b/app/src/main/kotlin/com/arflix/tv/data/repository/CloudSyncCoordinator.kt index 0f96711c1..fa87c29d2 100644 --- a/app/src/main/kotlin/com/arflix/tv/data/repository/CloudSyncCoordinator.kt +++ b/app/src/main/kotlin/com/arflix/tv/data/repository/CloudSyncCoordinator.kt @@ -8,6 +8,7 @@ import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.delay import kotlinx.coroutines.flow.collectLatest import kotlinx.coroutines.launch +import java.util.concurrent.atomic.AtomicBoolean import javax.inject.Inject import javax.inject.Singleton @@ -18,40 +19,48 @@ class CloudSyncCoordinator @Inject constructor( private val authRepository: AuthRepository ) { private val scope = CoroutineScope(SupervisorJob() + Dispatchers.IO) + private val lifecycleLock = Any() private var collectorJob: Job? = null private var flushJob: Job? = null - private val started = java.util.concurrent.atomic.AtomicBoolean(false) + private val started = AtomicBoolean(false) fun start() { - if (!started.compareAndSet(false, true)) return - collectorJob = scope.launch { - invalidationBus.events.collectLatest { invalidation -> - if (authRepository.getCurrentUserId().isNullOrBlank()) return@collectLatest - cloudSyncRepository.markLocalStateDirty() - scheduleFlush(invalidation) + synchronized(lifecycleLock) { + if (!started.compareAndSet(false, true)) return + collectorJob = scope.launch { + invalidationBus.events.collectLatest { invalidation -> + if (authRepository.getCurrentUserId().isNullOrBlank()) return@collectLatest + cloudSyncRepository.markLocalStateDirty() + scheduleFlush(invalidation) + } } } } fun stop() { - started.set(false) - collectorJob?.cancel() - flushJob?.cancel() - collectorJob = null - flushJob = null + synchronized(lifecycleLock) { + started.set(false) + collectorJob?.cancel() + flushJob?.cancel() + collectorJob = null + flushJob = null + } } private fun scheduleFlush(invalidation: CloudSyncInvalidation) { - flushJob?.cancel() - flushJob = scope.launch { - delay(debounceMsFor(invalidation.scope)) - if (authRepository.getCurrentUserId().isNullOrBlank()) return@launch - runCatching { cloudSyncRepository.pushToCloud() } - .onFailure { error -> - Log.w("CloudSyncCoordinator", "Cloud push failed after ${invalidation.scope}: ${error.message}") - cloudSyncRepository.markLocalStateDirty() - } + synchronized(lifecycleLock) { + if (!started.get()) return + flushJob?.cancel() + flushJob = scope.launch { + delay(debounceMsFor(invalidation.scope)) + if (authRepository.getCurrentUserId().isNullOrBlank()) return@launch + runCatching { cloudSyncRepository.pushToCloud() } + .onFailure { error -> + Log.w("CloudSyncCoordinator", "Cloud push failed after ${invalidation.scope}: ${error.message}") + cloudSyncRepository.markLocalStateDirty() + } + } } }