From e25e12df09584eb5544863a00d0e41cb130c5111 Mon Sep 17 00:00:00 2001 From: tapframe <85391825+tapframe@users.noreply.github.com> Date: Sun, 12 Jul 2026 19:30:31 +0530 Subject: [PATCH] fix: prevent settings sync reset regression --- .../app/core/sync/ProfileSettingsSync.kt | 316 ++---------------- 1 file changed, 32 insertions(+), 284 deletions(-) diff --git a/composeApp/src/commonMain/kotlin/com/nuvio/app/core/sync/ProfileSettingsSync.kt b/composeApp/src/commonMain/kotlin/com/nuvio/app/core/sync/ProfileSettingsSync.kt index 11fe09158..ad0774205 100644 --- a/composeApp/src/commonMain/kotlin/com/nuvio/app/core/sync/ProfileSettingsSync.kt +++ b/composeApp/src/commonMain/kotlin/com/nuvio/app/core/sync/ProfileSettingsSync.kt @@ -29,7 +29,6 @@ import com.nuvio.app.features.tmdb.TmdbSettingsStorage import com.nuvio.app.features.tmdb.TmdbSettingsRepository import com.nuvio.app.features.trakt.TraktCommentsStorage import com.nuvio.app.features.trakt.TraktCommentsSettings -import com.nuvio.app.features.trakt.ProfileSettingsWatchSourceOutbox import com.nuvio.app.features.trakt.TraktSettingsStorage import com.nuvio.app.features.trakt.TraktSettingsRepository import com.nuvio.app.features.watchprogress.ContinueWatchingPreferencesStorage @@ -37,21 +36,16 @@ import com.nuvio.app.features.watchprogress.ContinueWatchingPreferencesRepositor import io.github.jan.supabase.postgrest.postgrest import io.github.jan.supabase.postgrest.rpc import kotlin.concurrent.Volatile -import kotlinx.atomicfu.locks.SynchronizedObject -import kotlinx.atomicfu.locks.synchronized -import kotlinx.coroutines.CancellationException import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.FlowPreview import kotlinx.coroutines.Job import kotlinx.coroutines.SupervisorJob -import kotlinx.coroutines.delay import kotlinx.coroutines.flow.combine import kotlinx.coroutines.flow.debounce import kotlinx.coroutines.flow.distinctUntilChanged import kotlinx.coroutines.flow.drop import kotlinx.coroutines.flow.map -import kotlinx.coroutines.flow.onEach import kotlinx.coroutines.launch import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.withLock @@ -67,22 +61,10 @@ import kotlinx.serialization.json.put private const val PUSH_DEBOUNCE_MS = 1500L -private data class ObservedProfileSettingsChange( - val signature: String, - val accountId: String?, -) - -private data class SkippedProfileSettingsPush( - val signature: String, - val accountId: String?, - val profileId: Int, -) - object ProfileSettingsSync { private val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default) private val log = Logger.withTag("ProfileSettingsSync") private val syncMutex = Mutex() - private val observeLock = SynchronizedObject() private val json = Json { ignoreUnknownKeys = true encodeDefaults = true @@ -95,62 +77,33 @@ object ProfileSettingsSync { private var isServerSyncInFlight: Boolean = false @Volatile - private var skipNextPush: SkippedProfileSettingsPush? = null - - @Volatile - private var pushEnabledAccountId: String? = null - - @Volatile - private var pendingLocalPush: SkippedProfileSettingsPush? = null + private var skipNextPushSignature: String? = null private var observeJob: Job? = null - private var pendingPushRetryJob: Job? = null - fun startObserving() = synchronized(observeLock) { - if (observeJob?.isActive == true) return@synchronized + fun startObserving() { + if (observeJob?.isActive == true) return ensureRepositoriesLoaded() observeLocalChangesAndPush() } fun clearAccountState() { - synchronized(observeLock) { - observeJob?.cancel() - observeJob = null - } - skipNextPush = null - pushEnabledAccountId = null - pendingLocalPush = null - pendingPushRetryJob?.cancel() - pendingPushRetryJob = null + observeJob?.cancel() + observeJob = null + skipNextPushSignature = null } suspend fun pull(profileId: Int): Boolean { - startObserving() - val accountId = currentCloudAccountId() ?: return false + ensureRepositoriesLoaded() return syncMutex.withLock { - if (!isCurrentSyncTarget(profileId = profileId, accountId = accountId)) { + if (ProfileRepository.activeProfileId != profileId) { log.d { "pull(profileId=$profileId) — skipped because profile is no longer active" } return@withLock false } isServerSyncInFlight = true try { - pushEnabledAccountId = accountId - hydrateDurableWatchSourcePush(profileId = profileId, accountId = accountId) - if ( - !pushCurrentStateLocked( - profileId = profileId, - accountId = accountId, - forceCurrentState = false, - ) - ) { - schedulePendingPushRetry() - return@withLock false - } - val observedSignatureAtStart = currentObservedStateSignature() val localBlob = exportSettingsBlob() - if (!isCurrentSyncTarget(profileId = profileId, accountId = accountId)) { - throw CancellationException("Profile settings pull target changed") - } + if (ProfileRepository.activeProfileId != profileId) return@withLock false val localSignature = buildSignature(localBlob) val params = buildJsonObject { @@ -158,105 +111,41 @@ object ProfileSettingsSync { put("p_platform", MOBILE_SYNC_PLATFORM) } val result = SupabaseProvider.client.postgrest.rpc("sync_pull_profile_settings_blob", params) - if (!isCurrentSyncTarget(profileId = profileId, accountId = accountId)) { - throw CancellationException("Profile settings pull target changed") - } + if (ProfileRepository.activeProfileId != profileId) return@withLock false val response = result.decodeList().firstOrNull() val remoteJson = response?.settingsJson - val pendingDuringPull = pendingLocalPush?.let { pending -> - pending.accountId == accountId && pending.profileId == profileId - } == true - val durableSourceChangeDuringPull = - ProfileSettingsWatchSourceOutbox.pendingFor(accountId, profileId) != null - if ( - pendingDuringPull || - durableSourceChangeDuringPull || - currentObservedStateSignature() != observedSignatureAtStart - ) { - if ( - !pushCurrentStateLocked( - profileId = profileId, - accountId = accountId, - forceCurrentState = true, - ) - ) { - schedulePendingPushRetry() - } - return@withLock false - } - if (remoteJson == null) { log.i { "pull(profileId=$profileId) — no remote settings blob found" } - if (localSignature != defaultSignature()) { - pushToRemoteLocked(profileId, localBlob, accountId) - } - pushEnabledAccountId = accountId return@withLock false } - val remoteBlob = try { - json.decodeFromJsonElement(MobileProfileSettingsBlob.serializer(), remoteJson) - } catch (error: Throwable) { - log.e(error) { "pull(profileId=$profileId) — failed to decode remote settings blob" } - throw error - } - - var restoredPendingSourceAfterRemoteApply = false isApplyingRemoteBlob = true try { + val remoteBlob = runCatching { + json.decodeFromJsonElement(MobileProfileSettingsBlob.serializer(), remoteJson) + }.getOrElse { error -> + log.e(error) { "pull(profileId=$profileId) — failed to decode remote settings blob" } + return@withLock false + } val remoteSignature = buildSignature(remoteBlob) if (remoteSignature == localSignature) { log.d { "pull(profileId=$profileId) — remote matches local" } - pushEnabledAccountId = accountId return@withLock false } - if (!isCurrentSyncTarget(profileId = profileId, accountId = accountId)) { - throw CancellationException("Profile settings pull target changed") - } + if (ProfileRepository.activeProfileId != profileId) return@withLock false applyRemoteBlob(remoteBlob) - ProfileSettingsWatchSourceOutbox.pendingFor(accountId, profileId)?.let { pendingSource -> - if (TraktSettingsRepository.uiState.value.watchProgressSource != pendingSource.source) { - TraktSettingsRepository.setWatchProgressSource(pendingSource.source, profileId) - } - restoredPendingSourceAfterRemoteApply = true - } - skipNextPush = SkippedProfileSettingsPush( - signature = currentObservedStateSignature(), - accountId = currentCloudAccountId(), - profileId = profileId, - ) + skipNextPushSignature = currentObservedStateSignature() } finally { isApplyingRemoteBlob = false } - if (restoredPendingSourceAfterRemoteApply) { - pendingLocalPush = SkippedProfileSettingsPush( - signature = currentObservedStateSignature(), - accountId = accountId, - profileId = profileId, - ) - if ( - !pushCurrentStateLocked( - profileId = profileId, - accountId = accountId, - forceCurrentState = true, - ) - ) { - schedulePendingPushRetry() - } - return@withLock false - } - log.i { "pull(profileId=$profileId) — applied remote settings blob" } - pushEnabledAccountId = accountId true - } catch (error: CancellationException) { - throw error } catch (error: Exception) { log.e(error) { "pull(profileId=$profileId) — FAILED" } - throw error + false } finally { isServerSyncInFlight = false } @@ -265,22 +154,16 @@ object ProfileSettingsSync { suspend fun pushCurrentProfileToRemote(): Boolean { ensureRepositoriesLoaded() - val accountId = currentCloudAccountId() ?: return false return syncMutex.withLock { - try { + runCatching { val profileId = ProfileRepository.activeProfileId - if (!isCurrentSyncTarget(profileId = profileId, accountId = accountId)) return@withLock false - pushCurrentStateLocked( - profileId = profileId, - accountId = accountId, - forceCurrentState = true, - ) - } catch (error: CancellationException) { - throw error - } catch (error: Throwable) { + val blob = exportSettingsBlob() + if (ProfileRepository.activeProfileId != profileId) return@runCatching false + pushToRemoteLocked(profileId, blob) + true + }.onFailure { error -> log.e(error) { "pushCurrentProfileToRemote() — FAILED" } - false - } + }.getOrDefault(false) } } @@ -306,84 +189,24 @@ object ProfileSettingsSync { ) observeJob = scope.launch { - combine(signatureFlows) { - ObservedProfileSettingsChange( - signature = currentObservedStateSignature(), - accountId = currentCloudAccountId(), - ) - } + combine(signatureFlows) { currentObservedStateSignature() } .drop(1) .distinctUntilChanged() - .onEach { change -> - val authState = AuthRepository.state.value - if (authState !is AuthState.Authenticated || authState.isAnonymous) return@onEach - if (change.accountId == null || change.accountId != authState.userId) return@onEach - val observedChange = SkippedProfileSettingsPush( - signature = change.signature, - accountId = change.accountId, - profileId = ProfileRepository.activeProfileId, - ) - if (skipNextPush != observedChange) { - pendingLocalPush = observedChange - schedulePendingPushRetry() - } - } .debounce(PUSH_DEBOUNCE_MS) - .collect { change -> + .collect { signature -> val authState = AuthRepository.state.value if (authState !is AuthState.Authenticated || authState.isAnonymous) return@collect - if (change.accountId == null || change.accountId != authState.userId) return@collect - val profileId = ProfileRepository.activeProfileId - val observedChange = SkippedProfileSettingsPush( - signature = change.signature, - accountId = change.accountId, - profileId = profileId, - ) - if (skipNextPush == observedChange) { - skipNextPush = null - if (pendingLocalPush == observedChange) { - pendingLocalPush = null - } + if (isApplyingRemoteBlob || isServerSyncInFlight) return@collect + if (signature == skipNextPushSignature) { + skipNextPushSignature = null return@collect } - pendingLocalPush = observedChange - if (pushEnabledAccountId != change.accountId) return@collect - if (isApplyingRemoteBlob || isServerSyncInFlight) return@collect pushCurrentProfileToRemote() } } } - private fun schedulePendingPushRetry() { - if (pendingPushRetryJob?.isActive == true) return - pendingPushRetryJob = scope.launch { - var retryDelayMs = 5_000L - while (pendingLocalPush != null) { - delay(retryDelayMs) - val pending = pendingLocalPush ?: break - if ( - pending.accountId == currentCloudAccountId() && - pending.profileId == ProfileRepository.activeProfileId && - pushEnabledAccountId == pending.accountId && - !isApplyingRemoteBlob && - !isServerSyncInFlight && - pushCurrentProfileToRemote() - ) { - break - } - retryDelayMs = (retryDelayMs * 2L).coerceAtMost(60_000L) - } - } - } - - private suspend fun pushToRemoteLocked( - profileId: Int, - blob: MobileProfileSettingsBlob, - accountId: String, - ) { - if (!isCurrentSyncTarget(profileId = profileId, accountId = accountId)) { - throw CancellationException("Profile settings push target changed") - } + private suspend fun pushToRemoteLocked(profileId: Int, blob: MobileProfileSettingsBlob) { val params = buildJsonObject { put("p_profile_id", profileId) put("p_platform", MOBILE_SYNC_PLATFORM) @@ -391,74 +214,9 @@ object ProfileSettingsSync { putSyncOriginClientId() } SupabaseProvider.client.postgrest.rpc("sync_push_profile_settings_blob", params) - if (!isCurrentSyncTarget(profileId = profileId, accountId = accountId)) { - throw CancellationException("Profile settings push target changed") - } log.d { "pushToRemoteLocked(profileId=$profileId) — success" } } - private fun hydrateDurableWatchSourcePush(profileId: Int, accountId: String) { - val durableChange = ProfileSettingsWatchSourceOutbox.pendingFor(accountId, profileId) ?: return - if (TraktSettingsRepository.uiState.value.watchProgressSource != durableChange.source) { - TraktSettingsRepository.setWatchProgressSource(durableChange.source, profileId) - } - pendingLocalPush = SkippedProfileSettingsPush( - signature = currentObservedStateSignature(), - accountId = accountId, - profileId = profileId, - ) - schedulePendingPushRetry() - } - - private suspend fun pushCurrentStateLocked( - profileId: Int, - accountId: String, - forceCurrentState: Boolean, - ): Boolean { - val durableChange = ProfileSettingsWatchSourceOutbox.pendingFor(accountId, profileId) - val inMemoryChange = pendingLocalPush?.takeIf { pending -> - pending.accountId == accountId && pending.profileId == profileId - } - if (!forceCurrentState && durableChange == null && inMemoryChange == null) return true - - val signature = currentObservedStateSignature() - pushToRemoteLocked(profileId, exportSettingsBlob(), accountId) - if (currentObservedStateSignature() != signature) { - pendingLocalPush = SkippedProfileSettingsPush( - signature = currentObservedStateSignature(), - accountId = accountId, - profileId = profileId, - ) - schedulePendingPushRetry() - return false - } - - val pushedChange = SkippedProfileSettingsPush( - signature = signature, - accountId = accountId, - profileId = profileId, - ) - if (pendingLocalPush == pushedChange) { - pendingLocalPush = null - } - if ( - durableChange != null && - TraktSettingsRepository.uiState.value.watchProgressSource == durableChange.source - ) { - ProfileSettingsWatchSourceOutbox.clearIfMatches(durableChange) - } - - val durablePushRemains = ProfileSettingsWatchSourceOutbox.pendingFor(accountId, profileId) != null - val memoryPushRemains = pendingLocalPush?.let { pending -> - pending.accountId == accountId && pending.profileId == profileId - } == true - if (durablePushRemains || memoryPushRemains) { - schedulePendingPushRetry() - return false - } - return true - } - private fun exportSettingsBlob(): MobileProfileSettingsBlob { ensureRepositoriesLoaded() return MobileProfileSettingsBlob( @@ -547,9 +305,6 @@ object ProfileSettingsSync { private fun buildSignature(blob: MobileProfileSettingsBlob): String = json.encodeToString(MobileProfileSettingsBlob.serializer(), blob) - private fun defaultSignature(): String = - buildSignature(MobileProfileSettingsBlob()) - private fun currentObservedStateSignature(): String = listOf( "theme=${ThemeSettingsRepository.selectedTheme.value.name}", "amoled=${ThemeSettingsRepository.amoledEnabled.value}", @@ -569,13 +324,6 @@ object ProfileSettingsSync { "episode_release_alerts=${EpisodeReleaseNotificationsRepository.uiState.value.isEnabled}", ).joinToString(separator = "||") - private fun currentCloudAccountId(): String? = - (AuthRepository.state.value as? AuthState.Authenticated) - ?.takeUnless { it.isAnonymous } - ?.userId - - private fun isCurrentSyncTarget(profileId: Int, accountId: String): Boolean = - ProfileRepository.activeProfileId == profileId && currentCloudAccountId() == accountId } @Serializable