From 5f3aa487f5dbce69de7f73ce317d4e9404ebee49 Mon Sep 17 00:00:00 2001 From: tapframe <85391825+tapframe@users.noreply.github.com> Date: Wed, 18 Mar 2026 23:19:32 +0530 Subject: [PATCH] feat(trakt): episode remapping with concurrency and new mapping functions --- .../data/repository/TraktProgressService.kt | 161 +++++++++++++++++- 1 file changed, 159 insertions(+), 2 deletions(-) diff --git a/app/src/main/java/com/nuvio/tv/data/repository/TraktProgressService.kt b/app/src/main/java/com/nuvio/tv/data/repository/TraktProgressService.kt index dc98a07a..e416398e 100644 --- a/app/src/main/java/com/nuvio/tv/data/repository/TraktProgressService.kt +++ b/app/src/main/java/com/nuvio/tv/data/repository/TraktProgressService.kt @@ -27,7 +27,11 @@ import kotlinx.coroutines.CoroutineExceptionHandler import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.ExperimentalCoroutinesApi +import kotlinx.coroutines.Job import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.async +import kotlinx.coroutines.awaitAll +import kotlinx.coroutines.coroutineScope import kotlinx.coroutines.delay import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.collectLatest @@ -69,6 +73,7 @@ class TraktProgressService @Inject constructor( ) { companion object { private const val TAG = "TraktProgressSvc" + private const val REMOTE_EPISODE_REMAP_CONCURRENCY = 4 } private fun trace(message: String) { @@ -138,6 +143,7 @@ class TraktProgressService @Inject constructor( private val inFlightMetadataKeys = mutableSetOf() private val inFlightEpisodeProgressKeys = mutableSetOf() private val episodeProgressLastAttemptAtMs = mutableMapOf() + private var remoteRemapJob: Job? = null private var cachedMoviesPlayback: TimedCache>? = null private var cachedEpisodesPlayback: TimedCache>? = null private var cachedUserStats: TimedCache? = null @@ -157,6 +163,7 @@ class TraktProgressService @Inject constructor( @Volatile private var lastManualRefreshSignalMs: Long = 0L private val episodeProgressActivityVersion = AtomicLong(0L) + private val remoteSnapshotVersion = AtomicLong(0L) private val playbackCacheTtlMs = 30_000L private val userStatsCacheTtlMs = Long.MAX_VALUE @@ -648,10 +655,13 @@ class TraktProgressService @Inject constructor( } val snapshot = fetchAllProgressSnapshot(force = force) + val snapshotVersion = remoteSnapshotVersion.incrementAndGet() + remoteRemapJob?.cancel() remoteProgress.value = snapshot hasLoadedRemoteProgress.value = true reconcileOptimistic(snapshot) hydrateMetadata(snapshot) + scheduleRemoteEpisodeRemap(snapshotVersion, snapshot) } private suspend fun hasActivityChanged(): Boolean { @@ -877,7 +887,8 @@ class TraktProgressService @Inject constructor( toTraktUtcDateTime(System.currentTimeMillis() - windowMs) } val inProgressMovies = getPlayback("movies", force = force, startAt = playbackStartAt).mapNotNull { mapPlaybackMovie(it) } - val inProgressEpisodes = getPlayback("episodes", force = force, startAt = playbackStartAt).mapNotNull { mapPlaybackEpisode(it) } + val inProgressEpisodes = getPlayback("episodes", force = force, startAt = playbackStartAt) + .mapNotNull { mapPlaybackEpisodeRaw(it) } val mergedByKey = linkedMapOf() @@ -931,7 +942,7 @@ class TraktProgressService @Inject constructor( continue } - val mapped = mapEpisodeHistoryItem(item) ?: continue + val mapped = mapEpisodeHistoryItemRaw(item) ?: continue results.putIfAbsent(mapped.contentId, mapped) if (results.size >= maxRecentEpisodeHistoryEntries) { shouldStop = true @@ -947,6 +958,36 @@ class TraktProgressService @Inject constructor( return results.values.toList() } + private fun mapEpisodeHistoryItemRaw(item: TraktUserEpisodeHistoryItemDto): WatchProgress? { + val show = item.show ?: return null + val episode = item.episode ?: return null + val season = episode.season ?: return null + val number = episode.number ?: return null + + val contentId = normalizeContentId(show.ids) + if (contentId.isBlank()) return null + + return WatchProgress( + contentId = contentId, + contentType = "series", + name = show.title ?: contentId, + poster = null, + backdrop = null, + logo = null, + videoId = "$contentId:$season:$number", + season = season, + episode = number, + episodeTitle = episode.title, + position = 1L, + duration = 1L, + lastWatched = parseIsoToMillis(item.watchedAt), + progressPercent = 100f, + source = WatchProgress.SOURCE_TRAKT_HISTORY, + traktShowId = show.ids?.trakt, + traktEpisodeId = episode.ids?.trakt + ) + } + private suspend fun mapEpisodeHistoryItem(item: TraktUserEpisodeHistoryItemDto): WatchProgress? { val show = item.show ?: return null val episode = item.episode ?: return null @@ -1101,6 +1142,37 @@ class TraktProgressService @Inject constructor( ) } + private fun mapPlaybackEpisodeRaw(item: TraktPlaybackItemDto): WatchProgress? { + val show = item.show ?: return null + val episode = item.episode ?: return null + val season = episode.season ?: return null + val number = episode.number ?: return null + + val contentId = normalizeContentId(show.ids) + if (contentId.isBlank()) return null + + return WatchProgress( + contentId = contentId, + contentType = "series", + name = show.title ?: contentId, + poster = null, + backdrop = null, + logo = null, + videoId = "$contentId:$season:$number", + season = season, + episode = number, + episodeTitle = episode.title, + position = 0L, + duration = 0L, + lastWatched = parseIsoToMillis(item.pausedAt), + progressPercent = item.progress?.coerceIn(0f, 100f), + source = WatchProgress.SOURCE_TRAKT_PLAYBACK, + traktPlaybackId = item.id, + traktShowId = show.ids?.trakt, + traktEpisodeId = episode.ids?.trakt + ) + } + private suspend fun mapPlaybackEpisode(item: TraktPlaybackItemDto): WatchProgress? { val show = item.show ?: return null val episode = item.episode ?: return null @@ -1142,6 +1214,91 @@ class TraktProgressService @Inject constructor( ) } + private fun scheduleRemoteEpisodeRemap( + snapshotVersion: Long, + snapshot: List + ) { + if (snapshot.none(::shouldRemapEpisodeProgress)) return + + remoteRemapJob = scope.launch { + val remapped = runCatching { + remapEpisodeProgressSnapshot(snapshot) + }.getOrElse { error -> + Log.w(TAG, "remote episode remap failed", error) + return@launch + } + + if (snapshotVersion != remoteSnapshotVersion.get()) return@launch + if (remapped == snapshot) return@launch + + remoteProgress.value = remapped + reconcileOptimistic(remapped) + hydrateMetadata(remapped) + } + } + + private suspend fun remapEpisodeProgressSnapshot( + snapshot: List + ): List = coroutineScope { + val remapSemaphore = Semaphore(REMOTE_EPISODE_REMAP_CONCURRENCY) + snapshot.map { progress -> + async { + if (!shouldRemapEpisodeProgress(progress)) { + progress + } else { + remapSemaphore.withPermit { + remapEpisodeProgress(progress) + } + } + } + }.awaitAll() + } + + private fun shouldRemapEpisodeProgress(progress: WatchProgress): Boolean { + val normalizedType = progress.contentType.lowercase() + return normalizedType in listOf("series", "tv") && + progress.season != null && + progress.episode != null && + progress.source in setOf( + WatchProgress.SOURCE_TRAKT_HISTORY, + WatchProgress.SOURCE_TRAKT_PLAYBACK + ) + } + + private suspend fun remapEpisodeProgress(progress: WatchProgress): WatchProgress { + val season = progress.season ?: return progress + val episode = progress.episode ?: return progress + val resolvedEpisode = resolveAddonEpisodeProgress( + contentId = progress.contentId, + season = season, + episode = episode, + episodeTitle = progress.episodeTitle + ) ?: return progress.withResolvedEpisodeVideoId() + + val resolvedSeason = resolvedEpisode.season + val resolvedNumber = resolvedEpisode.episode + val videoId = resolvedEpisode.videoId + ?: resolveEpisodeVideoId(progress.contentId, resolvedSeason, resolvedNumber) + + return progress.copy( + videoId = videoId, + season = resolvedSeason, + episode = resolvedNumber, + episodeTitle = resolvedEpisode.title ?: progress.episodeTitle + ) + } + + private suspend fun WatchProgress.withResolvedEpisodeVideoId(): WatchProgress { + val season = season ?: return this + val episode = episode ?: return this + val placeholderVideoId = "${contentId}:$season:$episode" + if (videoId.isNotBlank() && videoId != placeholderVideoId) return this + + val resolvedVideoId = resolveEpisodeVideoId(contentId, season, episode) + if (resolvedVideoId == videoId) return this + return copy(videoId = resolvedVideoId) + } + private suspend fun mapSeasonProgress( contentId: String, season: TraktShowSeasonProgressDto