mirror of
https://github.com/tapframe/NuvioStreaming.git
synced 2026-08-04 10:36:56 +00:00
fix: prevent settings sync reset regression
This commit is contained in:
parent
d9954c04e8
commit
e25e12df09
1 changed files with 32 additions and 284 deletions
|
|
@ -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<SettingsBlobResponse>().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
|
||||
|
|
|
|||
Loading…
Reference in a new issue