feat(trakt): episode remapping with concurrency and new mapping functions

This commit is contained in:
tapframe 2026-03-18 23:19:32 +05:30
parent 6a4e1d373c
commit 5f3aa487f5

View file

@ -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<String>()
private val inFlightEpisodeProgressKeys = mutableSetOf<String>()
private val episodeProgressLastAttemptAtMs = mutableMapOf<String, Long>()
private var remoteRemapJob: Job? = null
private var cachedMoviesPlayback: TimedCache<List<TraktPlaybackItemDto>>? = null
private var cachedEpisodesPlayback: TimedCache<List<TraktPlaybackItemDto>>? = null
private var cachedUserStats: TimedCache<TraktCachedStats>? = 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<String, WatchProgress>()
@ -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<WatchProgress>
) {
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<WatchProgress>
): List<WatchProgress> = 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