From ab4c87676cc09626ec91bcff71cf54128b5616cd Mon Sep 17 00:00:00 2001 From: silentbil Date: Fri, 25 Sep 2026 16:28:48 +0300 Subject: [PATCH 1/3] telegram fixes --- .../tv/data/repository/CloudSyncRepository.kt | 10 +- .../tv/data/repository/StreamRepository.kt | 207 ++++++++-- .../arflix/tv/data/telegram/TelegramClient.kt | 5 + .../tv/data/telegram/TelegramRepository.kt | 46 ++- .../data/telegram/TelegramSourceResolver.kt | 359 ++++++++++++++---- .../arflix/tv/ui/components/StreamSelector.kt | 4 +- .../tv/ui/screens/details/DetailsScreen.kt | 36 +- .../tv/ui/screens/details/DetailsViewModel.kt | 86 ++++- .../tv/ui/screens/player/PlayerViewModel.kt | 18 + .../telegram/TelegramSettingsScreen.kt | 50 ++- .../telegram/TelegramSettingsViewModel.kt | 9 + app/src/main/res/values-iw/strings.xml | 5 + app/src/main/res/values/strings.xml | 5 + 13 files changed, 705 insertions(+), 135 deletions(-) diff --git a/app/src/main/kotlin/com/arflix/tv/data/repository/CloudSyncRepository.kt b/app/src/main/kotlin/com/arflix/tv/data/repository/CloudSyncRepository.kt index 825051973..c93df737d 100644 --- a/app/src/main/kotlin/com/arflix/tv/data/repository/CloudSyncRepository.kt +++ b/app/src/main/kotlin/com/arflix/tv/data/repository/CloudSyncRepository.kt @@ -469,7 +469,7 @@ class CloudSyncRepository @Inject constructor( "accentColor", "oledBlackBackground", "skipProfileSelection", "customUserAgent", "dnsProvider", "subtitleAiEnabled", "subtitleAiAutoSelect", "subtitleAiFindBestMatch", "subtitlePreloadEnabled", "dolbyVisionCompatEnabled", "subtitleAiApiKey", - "subtitleAiModel", "subtitleRemoveHearingImpaired" + "subtitleAiModel", "subtitleRemoveHearingImpaired", "telegramSearchOnClickOnly" ) // Per-profile fields excluded from the generic merge (handled by their own logic / not values). private val profileMergeExclude = setOf("defaultSubtitle", "subtitleSettingsUpdatedAt") @@ -750,6 +750,10 @@ class CloudSyncRepository @Inject constructor( root.put("subtitleAiApiKey", prefs[subtitleAiApiKeyKey] ?: "") root.put("subtitleAiModel", prefs[subtitleAiModelKey] ?: "GROQ_LLAMA_70B") root.put("subtitleRemoveHearingImpaired", prefs[subtitleRemoveHearingImpairedKey] ?: true) + root.put( + "telegramSearchOnClickOnly", + prefs[com.arflix.tv.data.telegram.TelegramRepository.KEY_SEARCH_ON_CLICK_ONLY] ?: true + ) root.put("activeProfileId", profileRepository.getActiveProfileId() ?: JSONObject.NULL) root.put("profiles", JSONArray(gson.toJson(profiles))) @@ -1747,6 +1751,10 @@ class CloudSyncRepository @Inject constructor( if (root.has("dolbyVisionCompatEnabled")) { prefs[dolbyVisionCompatKey] = root.optBoolean("dolbyVisionCompatEnabled", true) } + if (root.has("telegramSearchOnClickOnly")) { + prefs[com.arflix.tv.data.telegram.TelegramRepository.KEY_SEARCH_ON_CLICK_ONLY] = + root.optBoolean("telegramSearchOnClickOnly", true) + } val apiKey = root.optString("subtitleAiApiKey", "") if (apiKey.isNotBlank()) prefs[subtitleAiApiKeyKey] = apiKey val model = root.optString("subtitleAiModel", "GROQ_LLAMA_70B") diff --git a/app/src/main/kotlin/com/arflix/tv/data/repository/StreamRepository.kt b/app/src/main/kotlin/com/arflix/tv/data/repository/StreamRepository.kt index 0dfa00eb3..a10704c0d 100644 --- a/app/src/main/kotlin/com/arflix/tv/data/repository/StreamRepository.kt +++ b/app/src/main/kotlin/com/arflix/tv/data/repository/StreamRepository.kt @@ -1718,6 +1718,16 @@ class StreamRepository @Inject constructor( return "$profileId|$type|$imdbId|${season ?: 0}|${episode ?: 0}$providerPart|addons:$addonRevision" } + /** + * Where the addon-only lookups ([resolveMovieStreams], [resolveEpisodeStreams]) store their + * result. They share [streamCacheKey] with the progressive lookups, which also carry Telegram + * sources — and writing the addon-only list under that key replaced the complete one: The + * Shards S01E01 (Sept 2026) showed a Telegram source, the player refreshed its source list + * through the addon-only path, and reopening the episode within the cache's 10 minutes served + * the list without Telegram. They still READ the complete list first; they only never write it. + */ + private fun addonOnlyCacheKey(cacheKey: String): String = "$cacheKey|addon-only" + private fun cacheTtlMsFor(result: StreamResult): Long { val streams = result.streams if (streams.isEmpty()) return STREAM_RESULT_EMPTY_CACHE_TTL_MS @@ -2293,16 +2303,22 @@ class StreamRepository @Inject constructor( addonRevision = integrationCacheRevision(streamAddons) ) val profileId = profileManager.getProfileIdSync() + // This lookup is addon-only; see addonOnlyCacheKey for why it never writes the full key. + val ownCacheKey = addonOnlyCacheKey(cacheKey) if (!forceRefresh) { synchronized(streamResultCache) { - val cached = streamResultCache[cacheKey] - if (cached != null && isStreamCacheFresh(cached)) { - return@withContext cached.result + for (key in listOf(cacheKey, ownCacheKey)) { + val cached = streamResultCache[key] + if (cached != null && isStreamCacheFresh(cached)) { + return@withContext cached.result + } } } - loadPersistedStreamResult(profileId, cacheKey)?.let { persisted -> - synchronized(streamResultCache) { streamResultCache[cacheKey] = persisted } - return@withContext persisted.result + for (key in listOf(cacheKey, ownCacheKey)) { + loadPersistedStreamResult(profileId, key)?.let { persisted -> + synchronized(streamResultCache) { streamResultCache[key] = persisted } + return@withContext persisted.result + } } } @@ -2318,11 +2334,13 @@ class StreamRepository @Inject constructor( // IPTV VOD enrichment is appended separately in ViewModels. val result = StreamResult(streams = filteredStreams, subtitles = emptyList()) - val finalCacheKey = streamCacheKey( - profileId = profileId, - type = "movie", - imdbId = imdbId, - addonRevision = integrationCacheRevision(streamAddons) + val finalCacheKey = addonOnlyCacheKey( + streamCacheKey( + profileId = profileId, + type = "movie", + imdbId = imdbId, + addonRevision = integrationCacheRevision(streamAddons) + ) ) val cachedResult = CachedStreamResult(result, System.currentTimeMillis()) synchronized(streamResultCache) { @@ -2392,7 +2410,14 @@ class StreamRepository @Inject constructor( } val prioritizedAddons = prioritizeStreamingAddons(streamAddons) - val telegramEnabled = telegramSourceResolver.isEnabled() && streamIntegrationRepository.isIntegrationEnabled(StreamIntegrationType.TELEGRAM) + val telegramConnected = telegramSourceResolver.isEnabled() && streamIntegrationRepository.isIntegrationEnabled(StreamIntegrationType.TELEGRAM) + // "Search Telegram only on click" (Telegram settings): list only what an earlier search + // found; the details screen offers the search itself as a source row. + val telegramOnClick = telegramConnected && telegramSourceResolver.searchOnClickOnly() + val telegramEnabled = telegramConnected && !telegramOnClick + val telegramCached = if (telegramOnClick) { + telegramSourceResolver.cachedResults(title = title, imdbId = imdbId) + } else emptyList() if (prioritizedAddons.isEmpty() && !telegramEnabled) { Log.w( TAG, @@ -2429,8 +2454,12 @@ class StreamRepository @Inject constructor( ) val mutex = Mutex() - val aggregatedStreams = mutableListOf() + val aggregatedStreams = mutableListOf().apply { addAll(telegramCached) } var completed = 0 + // Set when the Telegram lookup did not run to the end (timed out, refused, or stopped + // because playback started). The list is then not the answer and is not cached — the + // next source list searches again (see TelegramResolution). + var telegramIncomplete = false val totalAddons = prioritizedAddons.size + (if (telegramEnabled) 1 else 0) suspend fun sendProgress() { @@ -2444,14 +2473,16 @@ class StreamRepository @Inject constructor( if (completed == totalAddons) { val createdAtMs = System.currentTimeMillis() val finalResult = StreamResult(filtered, emptyList()) - synchronized(streamResultCache) { - streamResultCache[cacheKey] = CachedStreamResult(finalResult, createdAtMs) + if (!telegramIncomplete) { + synchronized(streamResultCache) { + streamResultCache[cacheKey] = CachedStreamResult(finalResult, createdAtMs) + } + persistStreamResult( + profileId = profileId, + cacheKey = cacheKey, + cached = CachedStreamResult(finalResult, createdAtMs) + ) } - persistStreamResult( - profileId = profileId, - cacheKey = cacheKey, - cached = CachedStreamResult(finalResult, createdAtMs) - ) if (filtered.isEmpty()) { AppLogger.breadcrumb( tag = "Sources", @@ -2505,12 +2536,15 @@ class StreamRepository @Inject constructor( if (telegramEnabled) { val telegramStreams = try { - withTimeoutOrNull(3_500L) { - telegramSourceResolver.resolve(title = title, year = year, imdbId = imdbId, isMovie = true) - } ?: emptyList() + val resolution = withTimeoutOrNull(3_500L) { + telegramSourceResolver.resolveDetailed(title = title, year = year, imdbId = imdbId, isMovie = true) + } + if (resolution?.complete != true) mutex.withLock { telegramIncomplete = true } + resolution?.streams.orEmpty() } catch (e: Exception) { if (e is kotlinx.coroutines.CancellationException) throw e Log.e(TAG, "[StreamFetch][Movie] telegram resolve failed", e) + mutex.withLock { telegramIncomplete = true } emptyList() } val valid = telegramStreams.filter { stream -> @@ -2560,15 +2594,32 @@ class StreamRepository @Inject constructor( if (telegramEnabled) { launch { + // Telegram results are shown as they are found (the lookup runs in stages); + // each update is the whole list so far and replaces the previous one. + var shownTelegram: List = emptyList() val telegramStreams = try { - telegramSourceResolver.resolve(title = title, year = year, imdbId = imdbId, isMovie = true) + val resolution = telegramSourceResolver.resolveDetailed( + title = title, year = year, imdbId = imdbId, isMovie = true, + onFound = { partial -> + mutex.withLock { + aggregatedStreams.removeAll(shownTelegram) + aggregatedStreams.addAll(partial) + shownTelegram = partial + sendProgress() + } + } + ) + if (!resolution.complete) mutex.withLock { telegramIncomplete = true } + resolution.streams } catch (e: Exception) { if (e is kotlinx.coroutines.CancellationException) throw e Log.e(TAG, "[StreamFetch][Movie] telegram resolve failed", e) + mutex.withLock { telegramIncomplete = true } emptyList() } mutex.withLock { + aggregatedStreams.removeAll(shownTelegram) aggregatedStreams.addAll(telegramStreams) completed += 1 sendProgress() @@ -2637,6 +2688,44 @@ class StreamRepository @Inject constructor( }.orEmpty() } + /** + * Whether source lists should offer Telegram as a "Search Telegram" row instead of searching + * by themselves (the Telegram setting, when Telegram is connected and enabled). + */ + suspend fun isTelegramSearchOnClick(): Boolean = + telegramSourceResolver.isEnabled() && + streamIntegrationRepository.isIntegrationEnabled(StreamIntegrationType.TELEGRAM) && + telegramSourceResolver.searchOnClickOnly() + + /** + * The search the user asked for from a source list. Same lookup, cache key and staged results + * as the automatic one: [onFound] gets the whole list so far each time it grows. + */ + suspend fun searchTelegramNow( + mediaType: MediaType, + title: String, + year: Int?, + season: Int?, + episode: Int?, + imdbId: String, + onFound: suspend (List) -> Unit + ): List = withContext(Dispatchers.IO) { + if (mediaType == MediaType.MOVIE) { + telegramSourceResolver.resolveDetailed(title = title, year = year, imdbId = imdbId, isMovie = true, onFound = onFound) + } else { + telegramSourceResolver.resolveDetailed( + title = title, year = null, season = season, episode = episode, + imdbId = imdbId, isMovie = false, onFound = onFound + ) + }.streams + } + + /** Playback is starting: stop running Telegram searches and keep new ones quiet. */ + fun onPlaybackStarted() = telegramSourceResolver.onPlaybackStarted() + + /** The player closed; Telegram searches may notify again. */ + fun onPlaybackEnded() = telegramSourceResolver.onPlaybackEnded() + suspend fun resolveTelegramStreams( mediaType: MediaType, title: String, @@ -2649,6 +2738,15 @@ class StreamRepository @Inject constructor( if (!telegramSourceResolver.isEnabled() || !streamIntegrationRepository.isIntegrationEnabled(StreamIntegrationType.TELEGRAM)) { return@withContext emptyList() } + if (telegramSourceResolver.searchOnClickOnly()) { + // No search from the player: only what the user's own search already found. + return@withContext telegramSourceResolver.cachedResults( + title = title, + season = if (mediaType == MediaType.MOVIE) null else season ?: 1, + episode = if (mediaType == MediaType.MOVIE) null else episode ?: 1, + imdbId = imdbId.orEmpty() + ) + } withTimeoutOrNull(timeoutMs) { try { if (mediaType == MediaType.MOVIE) { @@ -2889,11 +2987,15 @@ class StreamRepository @Inject constructor( providerEpisodeId = animeQueryOverride, addonRevision = integrationCacheRevision(streamAddons) ) + // This lookup is addon-only; see addonOnlyCacheKey for why it never writes the full key. + val ownCacheKey = addonOnlyCacheKey(cacheKey) if (!forceRefresh) { synchronized(streamResultCache) { - val cached = streamResultCache[cacheKey] - if (cached != null && isStreamCacheFresh(cached)) { - return@withContext cached.result + for (key in listOf(cacheKey, ownCacheKey)) { + val cached = streamResultCache[key] + if (cached != null && isStreamCacheFresh(cached)) { + return@withContext cached.result + } } } } @@ -2922,7 +3024,7 @@ class StreamRepository @Inject constructor( val result = StreamResult(filteredStreams, subtitles) synchronized(streamResultCache) { - streamResultCache[cacheKey] = CachedStreamResult(result = result, createdAtMs = System.currentTimeMillis()) + streamResultCache[ownCacheKey] = CachedStreamResult(result = result, createdAtMs = System.currentTimeMillis()) } result } @@ -2992,7 +3094,14 @@ class StreamRepository @Inject constructor( } val prioritizedAddons = prioritizeStreamingAddons(streamAddons) - val telegramEnabled = telegramSourceResolver.isEnabled() && streamIntegrationRepository.isIntegrationEnabled(StreamIntegrationType.TELEGRAM) + val telegramConnected = telegramSourceResolver.isEnabled() && streamIntegrationRepository.isIntegrationEnabled(StreamIntegrationType.TELEGRAM) + // "Search Telegram only on click" (Telegram settings): list only what an earlier search + // found; the details screen offers the search itself as a source row. + val telegramOnClick = telegramConnected && telegramSourceResolver.searchOnClickOnly() + val telegramEnabled = telegramConnected && !telegramOnClick + val telegramCached = if (telegramOnClick) { + telegramSourceResolver.cachedResults(title = title, season = season, episode = episode, imdbId = imdbId) + } else emptyList() if (prioritizedAddons.isEmpty() && !telegramEnabled) { Log.w( TAG, @@ -3022,8 +3131,12 @@ class StreamRepository @Inject constructor( ) val mutex = Mutex() - val aggregatedStreams = mutableListOf() + val aggregatedStreams = mutableListOf().apply { addAll(telegramCached) } var completed = 0 + // Set when the Telegram lookup did not run to the end (timed out, refused, or stopped + // because playback started). The list is then not the answer and is not cached — the + // next source list searches again (see TelegramResolution). + var telegramIncomplete = false val totalAddons = prioritizedAddons.size + (if (telegramEnabled) 1 else 0) suspend fun sendProgress() { @@ -3037,8 +3150,10 @@ class StreamRepository @Inject constructor( if (completed == totalAddons) { val createdAtMs = System.currentTimeMillis() val finalResult = StreamResult(filtered, emptyList()) - synchronized(streamResultCache) { - streamResultCache[cacheKey] = CachedStreamResult(finalResult, createdAtMs) + if (!telegramIncomplete) { + synchronized(streamResultCache) { + streamResultCache[cacheKey] = CachedStreamResult(finalResult, createdAtMs) + } } if (filtered.isEmpty()) { AppLogger.breadcrumb( @@ -3105,8 +3220,8 @@ class StreamRepository @Inject constructor( if (telegramEnabled) { val telegramStreams = try { - withTimeoutOrNull(3_500L) { - telegramSourceResolver.resolve( + val resolution = withTimeoutOrNull(3_500L) { + telegramSourceResolver.resolveDetailed( title = title, year = null, season = season, @@ -3114,10 +3229,13 @@ class StreamRepository @Inject constructor( imdbId = imdbId, isMovie = false ) - } ?: emptyList() + } + if (resolution?.complete != true) mutex.withLock { telegramIncomplete = true } + resolution?.streams.orEmpty() } catch (e: Exception) { if (e is kotlinx.coroutines.CancellationException) throw e Log.e(TAG, "[StreamFetch][Episode] telegram resolve failed", e) + mutex.withLock { telegramIncomplete = true } emptyList() } val valid = telegramStreams.filter { stream -> @@ -3181,22 +3299,37 @@ class StreamRepository @Inject constructor( if (telegramEnabled) { launch { + // Telegram results are shown as they are found (the lookup runs in stages); + // each update is the whole list so far and replaces the previous one. + var shownTelegram: List = emptyList() val telegramStreams = try { - telegramSourceResolver.resolve( + val resolution = telegramSourceResolver.resolveDetailed( title = title, year = null, season = season, episode = episode, imdbId = imdbId, - isMovie = false + isMovie = false, + onFound = { partial -> + mutex.withLock { + aggregatedStreams.removeAll(shownTelegram) + aggregatedStreams.addAll(partial) + shownTelegram = partial + sendProgress() + } + } ) + if (!resolution.complete) mutex.withLock { telegramIncomplete = true } + resolution.streams } catch (e: Exception) { if (e is kotlinx.coroutines.CancellationException) throw e Log.e(TAG, "[StreamFetch][Episode] telegram resolve failed", e) + mutex.withLock { telegramIncomplete = true } emptyList() } mutex.withLock { + aggregatedStreams.removeAll(shownTelegram) aggregatedStreams.addAll(telegramStreams) completed += 1 sendProgress() diff --git a/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramClient.kt b/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramClient.kt index 6f8be73ef..7ad1f0ecb 100644 --- a/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramClient.kt +++ b/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramClient.kt @@ -81,6 +81,11 @@ class TelegramClient @Inject constructor( stepLog("library loaded OK") _authState.value = TelegramAuthState.Initializing try { + // TDLib's default verbosity writes every request, update and file event to + // logcat — thousands of lines a minute, enough for Android to start dropping the + // app's own logs ("chatty … expire"). Errors only. + runCatching { Client.execute(TdApi.SetLogVerbosityLevel(1)) } + .onFailure { Log.w(TAG, "SetLogVerbosityLevel failed: ${it.message}") } stepLog("calling Client.create") client = Client.create( { update -> handleUpdate(update) }, diff --git a/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramRepository.kt b/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramRepository.kt index 39bf9a8ed..2bed6f617 100644 --- a/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramRepository.kt +++ b/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramRepository.kt @@ -2,8 +2,12 @@ package com.arflix.tv.data.telegram import android.content.Context import android.util.Log +import androidx.datastore.preferences.core.booleanPreferencesKey import androidx.datastore.preferences.core.edit import androidx.datastore.preferences.core.stringPreferencesKey +import com.arflix.tv.data.repository.CloudSyncInvalidationBus +import com.arflix.tv.data.repository.CloudSyncScope +import com.arflix.tv.util.settingsDataStore import com.arflix.tv.util.telegramDataStore import dagger.hilt.android.qualifiers.ApplicationContext import kotlinx.coroutines.flow.Flow @@ -37,11 +41,18 @@ data class TelegramVideoMessage( class TelegramRepository @Inject constructor( @ApplicationContext private val context: Context, private val client: TelegramClient, - private val proxy: TelegramStreamingProxy + private val proxy: TelegramStreamingProxy, + private val invalidationBus: CloudSyncInvalidationBus ) { companion object { private const val TAG = "TelegramRepository" + private const val SEARCH_REQUEST_TIMEOUT_MS = 20_000L private val KEY_EXCLUDED_CHATS = stringPreferencesKey("excluded_chat_ids") + /** + * In the main settings store, not [telegramDataStore]: it is cloud-synced across devices + * (CloudSyncRepository "telegramSearchOnClickOnly"), and cloud sync only reads that store. + */ + val KEY_SEARCH_ON_CLICK_ONLY = booleanPreferencesKey("telegram_search_on_click_only") fun sessionMarker(context: Context) = File(context.filesDir, "tdlib_session_ok") @@ -180,28 +191,32 @@ class TelegramRepository @Inject constructor( } /** - * Searches globally across all chats (equivalent to Telethon's iter_messages(None, ...)). - * Runs two parallel searches — Document filter and Video filter — then merges results. + * Searches globally across all chats (equivalent to Telethon's iter_messages(None, ...)) and + * keeps the video files. + * + * ONE request, with no type filter; videos and video documents are picked out below. It used to + * be two per phrase (Document, then Video), and global search is exactly what Telegram rate- + * limits: one episode lookup sent 26 of them, most of which TDLib then held back until the + * caller gave up (Special Ops S1E1, Sept 2026: 8 of 9 requests unanswered after 10s). */ suspend fun searchVideoMessages( query: String, limit: Int = 50 ): List { - val filters = listOf( - TdApi.SearchMessagesFilterDocument(), - TdApi.SearchMessagesFilterVideo() - ) + val filters = listOf(null) val seen = mutableSetOf>() // dedupe by (fileName, fileSize) val results = mutableListOf() for (filter in filters) { + // Global search answers in 1–9s (Special Ops S1E1, Sept 2026: "פרק 1" 7.9s, + // "lioness s01e01" 8.9s); the default 10s cut those off on a slower second try. val result = client.sendRequest(TdApi.SearchMessages().also { req -> req.chatList = null // null = search all chats (like Telethon's iter_messages(None)) req.query = query req.offset = "" req.limit = limit req.filter = filter - }) + }, timeoutMs = SEARCH_REQUEST_TIMEOUT_MS) val found = (result as? TdApi.FoundMessages) ?: continue for (msg in found.messages) { @@ -248,6 +263,21 @@ class TelegramRepository @Inject constructor( fun getStreamUrl(fileId: Int): String = proxy.getUrl(fileId) + /** + * "Only search Telegram when clicking the Telegram source" (Telegram settings, on by default). + * Every source list used to search Telegram by itself — opening a show, pre-selecting an + * episode, the player's own list — and global search is what Telegram rate-limits, so the + * lookups nobody looked at used up the ones the user wanted. When on, a source list shows only + * what an earlier search found and offers the search as a row the user selects. + */ + val searchOnClickOnly: Flow = + context.settingsDataStore.data.map { prefs -> prefs[KEY_SEARCH_ON_CLICK_ONLY] ?: true } + + suspend fun setSearchOnClickOnly(enabled: Boolean) { + context.settingsDataStore.edit { prefs -> prefs[KEY_SEARCH_ON_CLICK_ONLY] = enabled } + invalidationBus.markDirty(CloudSyncScope.PROFILE_SETTINGS, reason = "telegram search on click") + } + fun getExcludedChatIds(): Flow> = context.telegramDataStore.data.map { prefs -> prefs[KEY_EXCLUDED_CHATS] diff --git a/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramSourceResolver.kt b/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramSourceResolver.kt index 1aa9436b6..56ec06a97 100644 --- a/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramSourceResolver.kt +++ b/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramSourceResolver.kt @@ -10,15 +10,32 @@ import com.arflix.tv.data.api.TmdbApi import com.arflix.tv.data.model.StreamSource import com.arflix.tv.util.Constants import dagger.hilt.android.qualifiers.ApplicationContext +import kotlinx.coroutines.CancellationException +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Deferred +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.async import kotlinx.coroutines.awaitAll import kotlinx.coroutines.coroutineScope +import kotlinx.coroutines.currentCoroutineContext +import kotlinx.coroutines.ensureActive +import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.first +import kotlinx.coroutines.launch import kotlinx.coroutines.withTimeoutOrNull import java.util.concurrent.ConcurrentHashMap +import java.util.concurrent.atomic.AtomicInteger import javax.inject.Inject import javax.inject.Singleton +/** + * A Telegram lookup's outcome. [complete] is false when the search did not run to the end — it + * timed out, Telegram refused it, or playback cancelled it — so [streams] says nothing about what + * the groups hold, and callers must not cache it as the answer. + */ +data class TelegramResolution(val streams: List, val complete: Boolean) + @Singleton class TelegramSourceResolver @Inject constructor( private val repository: TelegramRepository, @@ -29,8 +46,16 @@ class TelegramSourceResolver @Inject constructor( companion object { private const val TAG = "TelegramResolver" private const val SCORE_THRESHOLD = 55 - private const val SEARCH_TIMEOUT_MS = 20_000L + // The whole lookup. Results already found are kept (and shown) when it runs out; only a + // lookup that found nothing reports a timeout. + private const val SEARCH_TIMEOUT_MS = 30_000L private const val MAX_RESULTS = 100 + // Phrasings searched at once. Global search is what Telegram rate-limits, and TDLib then + // holds requests back rather than failing them: one episode lookup sent all ~13 phrasings + // (26 requests) together and most went unanswered until the lookup gave up (Special Ops + // S1E1, Sept 2026: 8 of 9 requests unanswered after 10s; limits of 4 and 8 in flight did + // not help — the volume did). + private const val MAX_PARALLEL_QUERIES = 2 private const val CACHE_TTL_SHORT_MS = 2 * 60 * 60 * 1_000L private const val CACHE_TTL_LONG_MS = 24 * 60 * 60 * 1_000L } @@ -38,12 +63,63 @@ class TelegramSourceResolver @Inject constructor( private data class CacheEntry(val results: List, val expiresAt: Long) private val cache = ConcurrentHashMap() + /** + * A running lookup: [found] holds every matching source so far (whole list, replaced as it + * grows) so callers can show results before the lookup ends; [result] is the final answer. + */ + private class Search( + val found: MutableStateFlow>, + val result: Deferred + ) + + /** + * One search per title/episode at a time: a second caller (the details screen and the player + * both ask) joins the running one instead of sending another burst. Runs in its own scope so + * [cancelActiveSearches] can stop it regardless of who started it. + */ + private val inFlight = ConcurrentHashMap() + private val searchScope = CoroutineScope(SupervisorJob() + Dispatchers.IO) + + /** Players currently open. While any is, searches stay silent — no toast over playback. */ + private val activePlayers = AtomicInteger(0) + private fun cacheKey(imdbId: String, title: String, season: Int?, episode: Int?) = if (imdbId.isNotBlank()) "$imdbId:${season ?: ""}:${episode ?: ""}" else "$title:${season ?: ""}:${episode ?: ""}" fun isEnabled(): Boolean = repository.isAuthenticated() + /** See [TelegramRepository.searchOnClickOnly]. */ + suspend fun searchOnClickOnly(): Boolean = repository.searchOnClickOnly.first() + + /** What an earlier completed lookup found for this title/episode, without searching. */ + fun cachedResults(title: String, season: Int? = null, episode: Int? = null, imdbId: String = ""): List = + cache[cacheKey(imdbId, title, season, episode)] + ?.takeIf { System.currentTimeMillis() < it.expiresAt } + ?.results + .orEmpty() + + /** + * Playback is starting: stop every running search (a source list the user has already chosen + * from no longer needs it) and keep later ones quiet until [onPlaybackEnded]. A stopped search + * caches nothing, so the next source list searches again. + */ + fun onPlaybackStarted() { + activePlayers.incrementAndGet() + cancelActiveSearches("playback started") + } + + fun onPlaybackEnded() { + activePlayers.updateAndGet { (it - 1).coerceAtLeast(0) } + } + + private fun cancelActiveSearches(reason: String) { + val running = inFlight.values.filter { it.result.isActive } + if (running.isEmpty()) return + Log.i(TAG, "Cancelling ${running.size} Telegram search(es): $reason") + running.forEach { it.result.cancel(CancellationException(reason)) } + } + // Old movies (released 2+ years ago) are stable — cache longer. // Series always use the short TTL since new episodes may appear at any time. private fun cacheTtl(year: Int?, isMovie: Boolean): Long { @@ -59,29 +135,110 @@ class TelegramSourceResolver @Inject constructor( episode: Int? = null, imdbId: String = "", isMovie: Boolean = true - ): List { - if (!repository.isAuthenticated()) return emptyList() + ): List = resolveDetailed(title, year, season, episode, imdbId, isMovie).streams + + /** + * [onFound] receives every matching source found so far, each time the list grows, while the + * lookup is still running — the complete list each time, to replace the previous one. + */ + suspend fun resolveDetailed( + title: String, + year: Int?, + season: Int? = null, + episode: Int? = null, + imdbId: String = "", + isMovie: Boolean = true, + onFound: (suspend (List) -> Unit)? = null + ): TelegramResolution { + if (!repository.isAuthenticated()) return TelegramResolution(emptyList(), complete = true) val key = cacheKey(imdbId, title, season, episode) cache[key]?.let { entry -> - if (System.currentTimeMillis() < entry.expiresAt) return entry.results + if (System.currentTimeMillis() < entry.expiresAt) return TelegramResolution(entry.results, complete = true) } - return try { - val results = withTimeoutOrNull(SEARCH_TIMEOUT_MS) { - resolveInternal(title, year, season, episode, imdbId, isMovie) - } ?: emptyList().also { - Log.w(TAG, "Telegram search timed out for '$title'") - showToast(context.getString(R.string.telegram_search_timed_out)) + val search = inFlight.compute(key) { _, running -> + running?.takeIf { it.result.isActive } ?: run { + val found = MutableStateFlow>(emptyList()) + Search(found, searchScope.async { + runSearch(key, title, year, season, episode, imdbId, isMovie, found) + }).also { started -> started.result.invokeOnCompletion { inFlight.remove(key, started) } } + } + }!! + return coroutineScope { + val relay = onFound?.let { callback -> + launch { search.found.collect { if (it.isNotEmpty()) callback(it) } } + } + try { + search.result.await() + } catch (e: CancellationException) { + // Either this caller was cancelled (rethrow) or the shared search was — by playback + // starting. The latter is not this caller's failure: keep what was found, cache nothing. + currentCoroutineContext().ensureActive() + TelegramResolution(search.found.value, complete = false) + } finally { + relay?.cancel() } + } + } + private suspend fun runSearch( + key: String, + title: String, + year: Int?, + season: Int?, + episode: Int?, + imdbId: String, + isMovie: Boolean, + found: MutableStateFlow> + ): TelegramResolution = try { + val startedAt = System.currentTimeMillis() + val finished = withTimeoutOrNull(SEARCH_TIMEOUT_MS) { + resolveInternal(title, year, season, episode, imdbId, isMovie, found) + } != null + val results = found.value + Log.i( + TAG, + "Telegram lookup '$title'${season?.let { " S${it}E$episode" }.orEmpty()}: " + + "${results.size} source(s)" + (if (finished) "" else ", timed out") + + " after ${System.currentTimeMillis() - startedAt}ms" + ) + if (!finished) { + // Not cached: a timeout says Telegram was slow, not that the groups have nothing. + // Caching it hid the sources for hours (2h for a series). What was found is kept. + if (results.isEmpty()) notifyUser(context.getString(R.string.telegram_search_timed_out)) + TelegramResolution(results, complete = false) + } else { cache[key] = CacheEntry(results, System.currentTimeMillis() + cacheTtl(year, isMovie)) - results - } catch (e: TelegramApiException) { - Log.w(TAG, "Telegram API error for '$title': ${e.message}") - showToast(friendlyError(e.message)) - emptyList() + TelegramResolution(results, complete = true) } + } catch (e: TelegramApiException) { + Log.w(TAG, "Telegram API error for '$title': ${e.message}") + if (found.value.isEmpty()) notifyUser(friendlyError(e.message)) + TelegramResolution(found.value, complete = false) + } + + /** + * Which phrasings run first. Measured on Special Ops S1E1 (Sept 2026), all 13 sent: ten files + * matched, and every one was found by `ע1פ1`, `ע1 פ1`, `פרק 1` or ` s01e01`; + * the other nine phrasings (`פ1`, `עונה 1 פרק 1`, the Hebrew title with Latin S/E patterns, + * the English title with `s1e1`/`s1 e1`/`s01 e01`) found nothing new and took 30 of the 43s. + * The core set always runs; the rest only when it found nothing. Movies build 2–6 phrasings + * and all of them are core. + */ + private fun splitCoreQueries(queries: List, season: Int?, episode: Int?): Pair, List> { + if (season == null || episode == null) return queries to emptyList() + val s = season.toString() + val e = episode.toString() + val sxe = "s${s.padStart(2, '0')}e${e.padStart(2, '0')}" + fun isCore(query: String): Boolean = + query.endsWith(" ע${s}פ$e") || + query.endsWith(" ע$s פ$e") || + (season == 1 && query.endsWith(" פרק $e") && !query.contains("עונה")) || + (season != 1 && query.endsWith(" עונה $s פרק $e")) || + (query.endsWith(" $sxe") && !matcher.isHebrew(query)) + val core = queries.filter(::isCore) + return if (core.isEmpty()) queries to emptyList() else core to queries.filterNot(::isCore) } private suspend fun resolveInternal( @@ -90,8 +247,9 @@ class TelegramSourceResolver @Inject constructor( season: Int?, episode: Int?, imdbId: String, - isMovie: Boolean - ): List { + isMovie: Boolean, + found: MutableStateFlow> + ) { val excludedIds = repository.getExcludedChatIds().first() // Read content language from SharedPreferences (same store SettingsViewModel writes to). // Avoids a DI cycle: StreamRepository → TelegramSourceResolver → MediaRepository → StreamRepository. @@ -109,80 +267,108 @@ class TelegramSourceResolver @Inject constructor( ) else matcher.buildMovieQueries(title, year, localizedTitle, englishTitle, originalTitle) + val (core, fallback) = splitCoreQueries(queries, season, episode) val seen = mutableSetOf>() - val allMessages = mutableListOf() - - coroutineScope { - queries.map { query -> - async { - try { - repository.searchVideoMessages(query, MAX_RESULTS) - .filter { it.chatId !in excludedIds } - } catch (e: TelegramApiException) { - throw e - } catch (e: kotlinx.coroutines.CancellationException) { - throw e - } catch (e: Exception) { - if (e is kotlinx.coroutines.CancellationException) throw e - - Log.e(TAG, "Search failed for '$query'", e) - emptyList() + // Matching sources by file, built once: the stream URL registers the file with the proxy. + val matched = LinkedHashMap, StreamSource>() + val searchStartedAt = System.currentTimeMillis() + + fun matchesRequest(msg: TelegramVideoMessage): Boolean = matcher.score( + fileName = msg.fileName, + caption = msg.caption, + title = title, + localizedTitle = localizedTitle, + englishTitle = englishTitle, + originalTitle = originalTitle, + year = year, + season = season, + episode = episode + ) >= SCORE_THRESHOLD + + suspend fun searchBatch(batch: List) { + val messages = coroutineScope { + batch.map { query -> + async { + try { + repository.searchVideoMessages(query, MAX_RESULTS) + .filter { it.chatId !in excludedIds } + } catch (e: TelegramApiException) { + throw e + } catch (e: kotlinx.coroutines.CancellationException) { + throw e + } catch (e: Exception) { + Log.e(TAG, "Search failed for '$query'", e) + emptyList() + } } - } - }.awaitAll().flatten().forEach { msg -> - if (seen.add(msg.fileName to msg.fileSize)) allMessages.add(msg) + }.awaitAll().flatten() } + var grew = false + for (msg in messages) { + val fileKey = msg.fileName to msg.fileSize + if (!seen.add(fileKey) || !matchesRequest(msg)) continue + matched[fileKey] = toStreamSource(msg) + grew = true + } + // Shown as soon as they are found; the lookup keeps going. + if (grew) found.value = sortSources(matched.values, langCode) } - return allMessages - .mapNotNull { msg -> - val score = matcher.score( - fileName = msg.fileName, - caption = msg.caption, - title = title, - localizedTitle = localizedTitle, - englishTitle = englishTitle, - originalTitle = originalTitle, - year = year, - season = season, - episode = episode - ) - if (score < SCORE_THRESHOLD) null else msg - } - .map { msg -> - val streamUrl = repository.getStreamUrl(msg.fileId) - val displayName = if (msg.fileName == "Default_Name.mkv" || msg.fileName == "Default_Name.mp4") - msg.caption.takeIf { it.isNotBlank() } ?: msg.fileName - else msg.fileName - val quality = parseQuality("${msg.fileName} ${msg.caption}") - StreamSource( - source = displayName, - addonName = "Telegram", - addonId = "telegram_native", - quality = quality, - size = formatBytes(msg.fileSize), - sizeBytes = msg.fileSize, - url = streamUrl, - infoHash = null, - fileIdx = null, - behaviorHints = com.arflix.tv.data.model.StreamBehaviorHints( - notWebReady = false, - filename = msg.fileName, - videoSize = msg.fileSize - ), - subtitles = emptyList(), - sources = emptyList(), - description = msg.caption.takeIf { it.isNotBlank() } - ) + var sent = 0 + for (batch in core.chunked(MAX_PARALLEL_QUERIES)) { + searchBatch(batch) + sent += batch.size + } + if (matched.isEmpty()) { + for (batch in fallback.chunked(MAX_PARALLEL_QUERIES)) { + searchBatch(batch) + sent += batch.size + if (matched.isNotEmpty()) break } - .sortedWith( - compareByDescending { langCode == "he" && matcher.isHebrew(it.source) } - .thenByDescending { qualityTier(it.quality) } - .thenByDescending { it.sizeBytes ?: 0L } - ) + } + Log.i( + TAG, + "Telegram search: $sent of ${queries.size} phrasings sent (${core.size} core, " + + "$MAX_PARALLEL_QUERIES at a time) in ${System.currentTimeMillis() - searchStartedAt}ms, " + + "${seen.size} distinct files, ${matched.size} matching" + ) } + private fun toStreamSource(msg: TelegramVideoMessage): StreamSource { + val streamUrl = repository.getStreamUrl(msg.fileId) + val displayName = if (msg.fileName == "Default_Name.mkv" || msg.fileName == "Default_Name.mp4") + msg.caption.takeIf { it.isNotBlank() } ?: msg.fileName + else msg.fileName + val quality = parseQuality("${msg.fileName} ${msg.caption}") + return StreamSource( + source = displayName, + addonName = "Telegram", + addonId = "telegram_native", + quality = quality, + size = formatBytes(msg.fileSize), + sizeBytes = msg.fileSize, + url = streamUrl, + infoHash = null, + fileIdx = null, + behaviorHints = com.arflix.tv.data.model.StreamBehaviorHints( + notWebReady = false, + filename = msg.fileName, + videoSize = msg.fileSize + ), + subtitles = emptyList(), + sources = emptyList(), + description = msg.caption.takeIf { it.isNotBlank() } + ) + } + + private fun sortSources(sources: Collection, langCode: String): List = + sources.sortedWith( + compareByDescending { langCode == "he" && matcher.isHebrew(it.source) } + .thenByDescending { qualityTier(it.quality) } + .thenByDescending { it.sizeBytes ?: 0L } + ) + private fun friendlyError(raw: String?): String { if (raw == null) return context.getString(R.string.telegram_search_failed) val waitSeconds = raw.removePrefix("FLOOD_WAIT_").toIntOrNull() @@ -192,6 +378,15 @@ class TelegramSourceResolver @Inject constructor( context.getString(R.string.telegram_error_raw, raw) } + /** A toast — unless a player is open: nothing about a background search belongs over playback. */ + private fun notifyUser(message: String) { + if (activePlayers.get() > 0) { + Log.i(TAG, "Not shown during playback: $message") + return + } + showToast(message) + } + private fun showToast(message: String) { Handler(Looper.getMainLooper()).post { Toast.makeText(context, message, Toast.LENGTH_LONG).show() diff --git a/app/src/main/kotlin/com/arflix/tv/ui/components/StreamSelector.kt b/app/src/main/kotlin/com/arflix/tv/ui/components/StreamSelector.kt index 30c5433c1..8bc5461c7 100644 --- a/app/src/main/kotlin/com/arflix/tv/ui/components/StreamSelector.kt +++ b/app/src/main/kotlin/com/arflix/tv/ui/components/StreamSelector.kt @@ -1199,7 +1199,9 @@ private fun sourceStatusText( val remaining = (totalAddons - completedAddons).coerceAtLeast(0) val elapsed = if (elapsedSeconds > 0 && (isLoading || pluginScrapersLoading)) "${elapsedSeconds}s \u2022 " else "" return when { - isLoading && totalAddons > 0 && remaining > 0 -> stringResource( + // Not gated on isLoading: once the first sources are listed isLoading is false, and a + // source still running (Telegram delivers in stages) read as "3/4 addons checked". + totalAddons > 0 && remaining > 0 -> stringResource( if (remaining == 1) { R.string.stream_status_still_checking_one } else { diff --git a/app/src/main/kotlin/com/arflix/tv/ui/screens/details/DetailsScreen.kt b/app/src/main/kotlin/com/arflix/tv/ui/screens/details/DetailsScreen.kt index 3f6d92675..1a0a19e8d 100644 --- a/app/src/main/kotlin/com/arflix/tv/ui/screens/details/DetailsScreen.kt +++ b/app/src/main/kotlin/com/arflix/tv/ui/screens/details/DetailsScreen.kt @@ -159,6 +159,7 @@ import com.arflix.tv.data.model.EpisodeIdentity import com.arflix.tv.data.model.MediaItem import com.arflix.tv.data.model.MediaType import com.arflix.tv.data.model.Review +import com.arflix.tv.data.model.StreamSource import com.arflix.tv.data.repository.MdbExternalRating import com.arflix.tv.network.OkHttpProvider import com.arflix.tv.ui.components.EpisodeContextMenu @@ -967,9 +968,23 @@ fun DetailsScreen( ) } // Stream Selector Modal + // With "only search Telegram when clicking", Telegram is offered as a row the user selects + // (see TelegramSearchRow). It is added here, for display only: autoplay, pre-warming and + // every other reader of uiState.streams never see it. + val telegramRowLabel = when (uiState.telegramSearchRow) { + TelegramSearchRow.HIDDEN -> null + TelegramSearchRow.IDLE -> stringResource(R.string.telegram_search_row_idle) + .takeIf { uiState.streams.none { it.addonId == TELEGRAM_SEARCH_ROW.addonId } } + TelegramSearchRow.SEARCHING -> stringResource(R.string.telegram_search_row_searching) + TelegramSearchRow.NONE -> stringResource(R.string.telegram_search_row_none) + } + val selectorStreams = remember(uiState.streams, telegramRowLabel) { + if (telegramRowLabel == null) uiState.streams + else uiState.streams + TELEGRAM_SEARCH_ROW.copy(source = telegramRowLabel) + } StreamSelector( isVisible = showStreamSelector, - streams = uiState.streams, + streams = selectorStreams, selectedStream = null, isLoading = uiState.isLoadingStreams, hasStreamingAddons = uiState.hasStreamingAddons, @@ -980,9 +995,13 @@ fun DetailsScreen( pluginScrapersLoading = uiState.pluginScrapersLoading, loadingPluginNames = uiState.loadingPluginNames, onFocusedStream = { stream -> - viewModel.prewarmStreamsAround(stream, uiState.streams) + if (!isTelegramSearchRow(stream)) viewModel.prewarmStreamsAround(stream, uiState.streams) }, onSelect = { stream -> + if (isTelegramSearchRow(stream)) { + viewModel.searchTelegramNow() + return@StreamSelector + } if (isPendingDebridStream(stream)) { viewModel.showToast( context.getString(R.string.details_toast_debrid_downloading), @@ -4850,3 +4869,16 @@ private fun SimilarMediaCard( onClick = onClick ) } + + +/** The "Search Telegram" row's stand-in source: listed under Telegram, never played. */ +private val TELEGRAM_SEARCH_ROW = StreamSource( + source = "", + addonName = "Telegram", + addonId = "telegram_native", + quality = "", + size = "", + url = "arvio-telegram-search://row" +) + +private fun isTelegramSearchRow(stream: StreamSource): Boolean = stream.url == TELEGRAM_SEARCH_ROW.url diff --git a/app/src/main/kotlin/com/arflix/tv/ui/screens/details/DetailsViewModel.kt b/app/src/main/kotlin/com/arflix/tv/ui/screens/details/DetailsViewModel.kt index d6ab695fe..0d52ecb2c 100644 --- a/app/src/main/kotlin/com/arflix/tv/ui/screens/details/DetailsViewModel.kt +++ b/app/src/main/kotlin/com/arflix/tv/ui/screens/details/DetailsViewModel.kt @@ -102,6 +102,8 @@ data class DetailsUiState( val loadingPluginNames: Set = emptySet(), val completedAddons: Int = 0, val totalAddons: Int = 0, + /** The "Search Telegram" source row (Telegram setting "only search when clicking"). */ + val telegramSearchRow: TelegramSearchRow = TelegramSearchRow.HIDDEN, val hasStreamingAddons: Boolean = true, val addonOrderedIds: List = emptyList(), val isInWatchlist: Boolean = false, @@ -211,7 +213,18 @@ enum class ToastType { } private fun isSupplementalStream(stream: StreamSource): Boolean = - IptvVodSourceIds.isIptvVodAddonId(stream.addonId) || stream.addonId == HomeServerRepository.ADDON_ID + IptvVodSourceIds.isIptvVodAddonId(stream.addonId) || stream.addonId == HomeServerRepository.ADDON_ID || + // Found by the user's own "Search Telegram" — outside the addon lookup, like the above. + stream.addonId == TELEGRAM_ADDON_ID + +private const val TELEGRAM_ADDON_ID = "telegram_native" + +/** + * The "Search Telegram" source row. With the Telegram setting "Only search Telegram when clicking + * the Telegram source" (on by default) source lists no longer search Telegram by themselves; the + * row starts the search, and what it finds is listed under it as it arrives. + */ +enum class TelegramSearchRow { HIDDEN, IDLE, SEARCHING, NONE } private fun Addon.isVodStreamingAddon(): Boolean = isEnabled && @@ -1899,6 +1912,58 @@ class DetailsViewModel @Inject constructor( } } + private data class TelegramSearchRequest( + val requestId: Long, + val mediaType: MediaType, + val title: String, + val year: Int?, + val season: Int?, + val episode: Int?, + val imdbId: String + ) + + /** What the "Search Telegram" row searches for — the current source list's title/episode. */ + private var telegramSearchRequest: TelegramSearchRequest? = null + + /** The user selected the "Search Telegram" row: search now, listing results as they arrive. */ + fun searchTelegramNow() { + val request = telegramSearchRequest ?: return + if (_uiState.value.telegramSearchRow == TelegramSearchRow.SEARCHING) return + _uiState.value = _uiState.value.copy(telegramSearchRow = TelegramSearchRow.SEARCHING) + viewModelScope.launch { + fun isCurrent() = request.requestId == loadStreamsRequestId + fun merge(found: List) { + if (!isCurrent() || found.isEmpty()) return + _uiState.value = _uiState.value.copy( + streams = sortPlayableStreamsFirst( + (_uiState.value.streams + found).distinctBy(::providerScopedStreamIdentity) + ), + isLoadingStreams = false + ) + } + val found = try { + streamRepository.searchTelegramNow( + mediaType = request.mediaType, + title = request.title, + year = request.year, + season = request.season, + episode = request.episode, + imdbId = request.imdbId, + onFound = { merge(it) } + ) + } catch (e: Exception) { + if (e is kotlinx.coroutines.CancellationException) throw e + Log.w(TAG, "[Telegram] search failed: ${e.message}") + emptyList() + } + if (!isCurrent()) return@launch + merge(found) + _uiState.value = _uiState.value.copy( + telegramSearchRow = if (found.isEmpty()) TelegramSearchRow.NONE else TelegramSearchRow.HIDDEN + ) + } + } + fun loadStreams(imdbId: String?, identity: EpisodeIdentity? = null) { val requestId = ++loadStreamsRequestId loadStreamsJob?.cancel() @@ -1918,8 +1983,10 @@ class DetailsViewModel @Inject constructor( streamsEpisodeIdentity = identity, subtitles = emptyList(), streamSearchStartTime = System.currentTimeMillis(), - pluginScrapersLoading = false + pluginScrapersLoading = false, + telegramSearchRow = TelegramSearchRow.HIDDEN ) + telegramSearchRequest = null val requestMediaType = currentMediaType val requestMediaId = currentMediaId @@ -1984,6 +2051,21 @@ class DetailsViewModel @Inject constructor( val originalLanguage = item?.originalLanguage val canonicalSeason = identity?.tmdbSeason val canonicalEpisode = identity?.tmdbEpisode + if (!effectiveStreamId.isNullOrBlank() && streamRepository.isTelegramSearchOnClick()) { + // The keys the automatic search would use, so both share its cache. + telegramSearchRequest = TelegramSearchRequest( + requestId = requestId, + mediaType = requestMediaType, + title = item?.title.orEmpty(), + year = item?.year?.toIntOrNull(), + season = if (requestMediaType == MediaType.MOVIE) null else canonicalSeason ?: 1, + episode = if (requestMediaType == MediaType.MOVIE) null else canonicalEpisode ?: 1, + imdbId = effectiveStreamId + ) + if (isCurrentRequest()) { + _uiState.value = _uiState.value.copy(telegramSearchRow = TelegramSearchRow.IDLE) + } + } val animeQueryOverride = identity?.kitsuQuery val homeServerEnabled = streamIntegrationRepository.isIntegrationEnabled(StreamIntegrationType.HOME_SERVER) val hasHomeServerConnections = homeServerEnabled && streamRepository.hasHomeServerConnections() diff --git a/app/src/main/kotlin/com/arflix/tv/ui/screens/player/PlayerViewModel.kt b/app/src/main/kotlin/com/arflix/tv/ui/screens/player/PlayerViewModel.kt index 7145dc071..c2a4e6950 100644 --- a/app/src/main/kotlin/com/arflix/tv/ui/screens/player/PlayerViewModel.kt +++ b/app/src/main/kotlin/com/arflix/tv/ui/screens/player/PlayerViewModel.kt @@ -646,6 +646,17 @@ class PlayerViewModel @Inject constructor( "has_selected_url" to (!_uiState.value.selectedStreamUrl.isNullOrBlank()).toString() ).apply { putAll(extra) } + /** This player has told the Telegram search that playback is on (see loadMedia / onCleared). */ + private var holdsTelegramQuiet = false + + override fun onCleared() { + if (holdsTelegramQuiet) { + holdsTelegramQuiet = false + streamRepository.onPlaybackEnded() + } + super.onCleared() + } + fun loadMedia( mediaType: MediaType, mediaId: Int, @@ -665,6 +676,13 @@ class PlayerViewModel @Inject constructor( airDate: String? = null ) { currentAirDate = airDate + // Playback is starting: Telegram searches still running for the source list the user + // just chose from are stopped, and none may toast over the player (they used to: "Telegram + // search timed out" appearing mid-episode). Held until this player is cleared. + if (!holdsTelegramQuiet) { + holdsTelegramQuiet = true + streamRepository.onPlaybackStarted() + } currentMediaType = mediaType currentMediaId = mediaId currentSeason = seasonNumber diff --git a/app/src/main/kotlin/com/arflix/tv/ui/screens/settings/telegram/TelegramSettingsScreen.kt b/app/src/main/kotlin/com/arflix/tv/ui/screens/settings/telegram/TelegramSettingsScreen.kt index 94242ee52..ab797f543 100644 --- a/app/src/main/kotlin/com/arflix/tv/ui/screens/settings/telegram/TelegramSettingsScreen.kt +++ b/app/src/main/kotlin/com/arflix/tv/ui/screens/settings/telegram/TelegramSettingsScreen.kt @@ -88,6 +88,7 @@ fun TelegramSettingsScreen( ) { val authState by viewModel.authState.collectAsState() val cacheSizeBytes by viewModel.cacheSizeBytes.collectAsState() + val searchOnClickOnly by viewModel.searchOnClickOnly.collectAsState() var showDisconnectConfirm by remember { mutableStateOf(false) } val isMobile = LocalDeviceType.current.isTouchDevice() val backMotion = rememberArvioPredictiveBack(enabled = isMobile && showHeader && !showDisconnectConfirm) { @@ -157,8 +158,10 @@ fun TelegramSettingsScreen( is TelegramAuthState.Ready -> ConnectedContent( firstName = state.firstName, cacheSizeBytes = cacheSizeBytes, + searchOnClickOnly = searchOnClickOnly, onDisconnect = { showDisconnectConfirm = true }, - onClearCache = { viewModel.clearCache() } + onClearCache = { viewModel.clearCache() }, + onToggleSearchOnClickOnly = { viewModel.setSearchOnClickOnly(!searchOnClickOnly) } ) is TelegramAuthState.Error -> ErrorContent( message = state.message, @@ -493,8 +496,10 @@ private fun PasswordContent(onSubmit: (String) -> Unit) { private fun ConnectedContent( firstName: String, cacheSizeBytes: Long, + searchOnClickOnly: Boolean, onDisconnect: () -> Unit, - onClearCache: () -> Unit + onClearCache: () -> Unit, + onToggleSearchOnClickOnly: () -> Unit ) { Column( modifier = Modifier.fillMaxWidth(), @@ -577,6 +582,47 @@ private fun ConnectedContent( ) } } + + Spacer(modifier = Modifier.height(12.dp)) + + var searchModeFocused by remember { mutableStateOf(false) } + Row( + modifier = Modifier + .fillMaxWidth() + .clickable { onToggleSearchOnClickOnly() } + .onFocusChanged { searchModeFocused = it.isFocused } + .background( + if (searchModeFocused) Pink.copy(alpha = 0.18f) else Color.White.copy(alpha = 0.04f), + RoundedCornerShape(12.dp) + ) + .border( + width = if (searchModeFocused) 2.dp else 1.dp, + color = if (searchModeFocused) Pink else Color.White.copy(alpha = 0.08f), + shape = RoundedCornerShape(12.dp) + ) + .padding(horizontal = 16.dp, vertical = 14.dp), + verticalAlignment = Alignment.CenterVertically, + horizontalArrangement = Arrangement.SpaceBetween + ) { + Column(modifier = Modifier.weight(1f)) { + Text( + text = stringResource(R.string.telegram_search_on_click_title), + style = ArflixTypography.cardTitle.copy(fontSize = 14.sp), + color = if (searchModeFocused) Pink else TextPrimary + ) + Text( + text = stringResource(R.string.telegram_search_on_click_desc), + style = ArflixTypography.caption.copy(fontSize = 13.sp), + color = TextSecondary + ) + } + Spacer(modifier = Modifier.width(12.dp)) + Text( + text = stringResource(if (searchOnClickOnly) R.string.on else R.string.off), + style = ArflixTypography.label.copy(fontSize = 11.sp), + color = if (searchOnClickOnly) SuccessGreen else TextSecondary + ) + } } } diff --git a/app/src/main/kotlin/com/arflix/tv/ui/screens/settings/telegram/TelegramSettingsViewModel.kt b/app/src/main/kotlin/com/arflix/tv/ui/screens/settings/telegram/TelegramSettingsViewModel.kt index f18fecb50..ea702762f 100644 --- a/app/src/main/kotlin/com/arflix/tv/ui/screens/settings/telegram/TelegramSettingsViewModel.kt +++ b/app/src/main/kotlin/com/arflix/tv/ui/screens/settings/telegram/TelegramSettingsViewModel.kt @@ -7,8 +7,10 @@ import com.arflix.tv.data.telegram.TelegramRepository import dagger.hilt.android.lifecycle.HiltViewModel import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.flow.MutableStateFlow +import kotlinx.coroutines.flow.SharingStarted import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.asStateFlow +import kotlinx.coroutines.flow.stateIn import kotlinx.coroutines.launch import javax.inject.Inject @@ -19,6 +21,13 @@ class TelegramSettingsViewModel @Inject constructor( val authState: StateFlow = repository.authState + val searchOnClickOnly: StateFlow = repository.searchOnClickOnly + .stateIn(viewModelScope, SharingStarted.WhileSubscribed(5_000L), true) + + fun setSearchOnClickOnly(enabled: Boolean) { + viewModelScope.launch { repository.setSearchOnClickOnly(enabled) } + } + private val _cacheSizeBytes = MutableStateFlow(0L) val cacheSizeBytes: StateFlow = _cacheSizeBytes.asStateFlow() diff --git a/app/src/main/res/values-iw/strings.xml b/app/src/main/res/values-iw/strings.xml index 5f2c7f626..9ad28a604 100644 --- a/app/src/main/res/values-iw/strings.xml +++ b/app/src/main/res/values-iw/strings.xml @@ -250,4 +250,9 @@ הצג את שורות האוספים המובנות של ARVIO (שירותים / ז׳אנרים / זיכיונות). כבה כדי להשתמש רק באוספים שייבאת אוספים הותקנו: %1$s (%2$d שורות) זה לא קובץ אוספים תקין (נדרש JSON של ייצוא אוספים מ-Nuvio) + חפש בטלגרם רק בלחיצה על מקור הטלגרם + מונע חריגה ממגבלת החיפוש של טלגרם. כבה כדי לחפש בטלגרם אוטומטית יחד עם שאר המקורות. + לחץ לחיפוש בטלגרם + מחפש בטלגרם… + לא נמצאו מקורות בטלגרם — לחץ לחיפוש חוזר diff --git a/app/src/main/res/values/strings.xml b/app/src/main/res/values/strings.xml index cba8e7084..9664f41d2 100644 --- a/app/src/main/res/values/strings.xml +++ b/app/src/main/res/values/strings.xml @@ -960,6 +960,11 @@ DISCONNECT Video Cache CLEAR + Only search Telegram when clicking the Telegram source + Avoids Telegram\'s search limit. Turn off to search Telegram automatically with the other sources. + Click to search Telegram + Searching Telegram… + No Telegram sources found — select to search again Connection Failed TRY AGAIN Disconnect Telegram? From 4bde8812caedab98d037bd25320e22ccad341b48 Mon Sep 17 00:00:00 2001 From: silentbil Date: Sat, 26 Sep 2026 14:53:03 +0300 Subject: [PATCH 2/3] =?UTF-8?q?fix(telegram):=20address=20review=20?= =?UTF-8?q?=E2=80=94=20unanswered=20requests,=20cut-off=20pages,=20reopene?= =?UTF-8?q?d=20Sources?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - A request Telegram did not answer in time now makes the whole lookup incomplete, so it is never cached as "no results" (the per-request 20s timeout used to look like an empty answer). - When a phrasing's unfiltered first page has more results behind it, the phrasing is repeated with the video and document filters, so newer text or photo posts cannot push the files off the only page read. Pages with nothing behind them still cost one request. - Telegram results found by a click-search are merged into every saved source list that Sources returns (fresh, stale, disk and addon-less paths), so reopening Sources keeps them, including with Telegram as the only provider. - The phrase loop and the result cache moved into TelegramPhraseSearch / TelegramResultCache (no Android dependencies) with regression tests for a per-request timeout, a cut-off first page, and reopening after a click-search (including Telegram-only). --- .../tv/data/repository/StreamRepository.kt | 66 ++++--- .../tv/data/telegram/TelegramPhraseSearch.kt | 128 ++++++++++++ .../tv/data/telegram/TelegramRepository.kt | 115 +++++------ .../data/telegram/TelegramSourceResolver.kt | 126 +++++------- .../repository/TelegramCachedSourcesTest.kt | 50 +++++ .../data/telegram/TelegramPhraseSearchTest.kt | 185 ++++++++++++++++++ 6 files changed, 502 insertions(+), 168 deletions(-) create mode 100644 app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramPhraseSearch.kt create mode 100644 app/src/test/kotlin/com/arflix/tv/data/repository/TelegramCachedSourcesTest.kt create mode 100644 app/src/test/kotlin/com/arflix/tv/data/telegram/TelegramPhraseSearchTest.kt diff --git a/app/src/main/kotlin/com/arflix/tv/data/repository/StreamRepository.kt b/app/src/main/kotlin/com/arflix/tv/data/repository/StreamRepository.kt index a10704c0d..608ba0c49 100644 --- a/app/src/main/kotlin/com/arflix/tv/data/repository/StreamRepository.kt +++ b/app/src/main/kotlin/com/arflix/tv/data/repository/StreamRepository.kt @@ -150,6 +150,18 @@ internal fun shouldTryNativeAnimeFallback( return language.isNullOrBlank() || language == "ja" } +/** + * A source list with the Telegram results an earlier Telegram search found. With "only search + * Telegram when clicking" those live in the Telegram lookup's cache, not in the saved source list, + * so every path that returns a saved list adds them here. + */ +internal fun withCachedTelegramSources( + streams: List, + telegramCached: List +): List = + if (telegramCached.isEmpty()) streams + else (streams + telegramCached).distinctBy(::providerScopedStreamIdentity) + internal fun providerScopedStreamIdentity(stream: StreamSource): String { return listOf( stream.addonId.trim(), @@ -2371,13 +2383,23 @@ class StreamRepository @Inject constructor( addonRevision = integrationCacheRevision(streamAddons) ) val cacheKey = if (sequential) "$baseCacheKey:seq" else baseCacheKey + val telegramConnected = telegramSourceResolver.isEnabled() && streamIntegrationRepository.isIntegrationEnabled(StreamIntegrationType.TELEGRAM) + // "Search Telegram only on click" (Telegram settings): list only what an earlier search + // found; the details screen offers the search itself as a source row. Worked out before + // the saved-list checks below: every list returned — saved, stale or addon-less — carries + // what a click-search already found, or reopening Sources dropped it (PR #757 review). + val telegramOnClick = telegramConnected && telegramSourceResolver.searchOnClickOnly() + val telegramEnabled = telegramConnected && !telegramOnClick + val telegramCached = if (telegramOnClick) { + telegramSourceResolver.cachedResults(title = title, imdbId = imdbId) + } else emptyList() if (!forceRefresh) { var warmCache: CachedStreamResult? = null synchronized(streamResultCache) { val cached = streamResultCache[cacheKey] if (cached != null) { if (isStreamCacheFresh(cached)) { - trySend(ProgressiveStreamResult(cached.result.streams, cached.result.subtitles, 1, 1, true)) + trySend(ProgressiveStreamResult(withCachedTelegramSources(cached.result.streams, telegramCached), cached.result.subtitles, 1, 1, true)) close() return@launch } @@ -2395,7 +2417,7 @@ class StreamRepository @Inject constructor( synchronized(streamResultCache) { streamResultCache[cacheKey] = cached } trySend( ProgressiveStreamResult( - streams = cached.result.streams, + streams = withCachedTelegramSources(cached.result.streams, telegramCached), subtitles = cached.result.subtitles, completedAddons = 1, totalAddons = 1, @@ -2410,14 +2432,6 @@ class StreamRepository @Inject constructor( } val prioritizedAddons = prioritizeStreamingAddons(streamAddons) - val telegramConnected = telegramSourceResolver.isEnabled() && streamIntegrationRepository.isIntegrationEnabled(StreamIntegrationType.TELEGRAM) - // "Search Telegram only on click" (Telegram settings): list only what an earlier search - // found; the details screen offers the search itself as a source row. - val telegramOnClick = telegramConnected && telegramSourceResolver.searchOnClickOnly() - val telegramEnabled = telegramConnected && !telegramOnClick - val telegramCached = if (telegramOnClick) { - telegramSourceResolver.cachedResults(title = title, imdbId = imdbId) - } else emptyList() if (prioritizedAddons.isEmpty() && !telegramEnabled) { Log.w( TAG, @@ -2431,19 +2445,19 @@ class StreamRepository @Inject constructor( if (!forceRefresh) { val cached = synchronized(streamResultCache) { streamResultCache[cacheKey] } if (cached != null) { - trySend(ProgressiveStreamResult(cached.result.streams, cached.result.subtitles, 1, 1, true)) + trySend(ProgressiveStreamResult(withCachedTelegramSources(cached.result.streams, telegramCached), cached.result.subtitles, 1, 1, true)) close() return@launch } val persisted = loadPersistedStreamResult(profileId = profileId, cacheKey = cacheKey) if (persisted != null) { synchronized(streamResultCache) { streamResultCache[cacheKey] = persisted } - trySend(ProgressiveStreamResult(persisted.result.streams, persisted.result.subtitles, 1, 1, true)) + trySend(ProgressiveStreamResult(withCachedTelegramSources(persisted.result.streams, telegramCached), persisted.result.subtitles, 1, 1, true)) close() return@launch } } - trySend(ProgressiveStreamResult(emptyList(), emptyList(), 0, 0, true)) + trySend(ProgressiveStreamResult(telegramCached, emptyList(), 0, 0, true)) close() return@launch } @@ -3067,13 +3081,23 @@ class StreamRepository @Inject constructor( addonRevision = integrationCacheRevision(streamAddons) ) val cacheKey = if (sequential) "$baseCacheKey:seq" else baseCacheKey + val telegramConnected = telegramSourceResolver.isEnabled() && streamIntegrationRepository.isIntegrationEnabled(StreamIntegrationType.TELEGRAM) + // "Search Telegram only on click" (Telegram settings): list only what an earlier search + // found; the details screen offers the search itself as a source row. Worked out before + // the saved-list checks below: every list returned — saved, stale or addon-less — carries + // what a click-search already found, or reopening Sources dropped it (PR #757 review). + val telegramOnClick = telegramConnected && telegramSourceResolver.searchOnClickOnly() + val telegramEnabled = telegramConnected && !telegramOnClick + val telegramCached = if (telegramOnClick) { + telegramSourceResolver.cachedResults(title = title, season = season, episode = episode, imdbId = imdbId) + } else emptyList() if (!forceRefresh) { var staleCache: CachedStreamResult? = null synchronized(streamResultCache) { val cached = streamResultCache[cacheKey] if (cached != null) { if (isStreamCacheFresh(cached)) { - trySend(ProgressiveStreamResult(cached.result.streams, cached.result.subtitles, 1, 1, true)) + trySend(ProgressiveStreamResult(withCachedTelegramSources(cached.result.streams, telegramCached), cached.result.subtitles, 1, 1, true)) close() return@launch } @@ -3083,7 +3107,7 @@ class StreamRepository @Inject constructor( staleCache?.let { cached -> trySend( ProgressiveStreamResult( - streams = cached.result.streams, + streams = withCachedTelegramSources(cached.result.streams, telegramCached), subtitles = cached.result.subtitles, completedAddons = 0, totalAddons = 1, @@ -3094,14 +3118,6 @@ class StreamRepository @Inject constructor( } val prioritizedAddons = prioritizeStreamingAddons(streamAddons) - val telegramConnected = telegramSourceResolver.isEnabled() && streamIntegrationRepository.isIntegrationEnabled(StreamIntegrationType.TELEGRAM) - // "Search Telegram only on click" (Telegram settings): list only what an earlier search - // found; the details screen offers the search itself as a source row. - val telegramOnClick = telegramConnected && telegramSourceResolver.searchOnClickOnly() - val telegramEnabled = telegramConnected && !telegramOnClick - val telegramCached = if (telegramOnClick) { - telegramSourceResolver.cachedResults(title = title, season = season, episode = episode, imdbId = imdbId) - } else emptyList() if (prioritizedAddons.isEmpty() && !telegramEnabled) { Log.w( TAG, @@ -3115,12 +3131,12 @@ class StreamRepository @Inject constructor( if (!forceRefresh) { val cached = synchronized(streamResultCache) { streamResultCache[cacheKey] } if (cached != null) { - trySend(ProgressiveStreamResult(cached.result.streams, cached.result.subtitles, 1, 1, true)) + trySend(ProgressiveStreamResult(withCachedTelegramSources(cached.result.streams, telegramCached), cached.result.subtitles, 1, 1, true)) close() return@launch } } - trySend(ProgressiveStreamResult(emptyList(), emptyList(), 0, 0, true)) + trySend(ProgressiveStreamResult(telegramCached, emptyList(), 0, 0, true)) close() return@launch } diff --git a/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramPhraseSearch.kt b/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramPhraseSearch.kt new file mode 100644 index 000000000..e5d719a4e --- /dev/null +++ b/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramPhraseSearch.kt @@ -0,0 +1,128 @@ +package com.arflix.tv.data.telegram + +import kotlinx.coroutines.async +import kotlinx.coroutines.awaitAll +import kotlinx.coroutines.coroutineScope + +/** Which messages a global search page asks Telegram for. */ +enum class TelegramSearchFilter { ALL, VIDEO, DOCUMENT } + +/** + * One page of a global Telegram search. + * + * [answered] is false when Telegram gave no answer in time (TDLib holds rate-limited requests + * rather than failing them) — the page says nothing about what the groups hold. [hasMore] is true + * when the search has results beyond this page. + */ +data class TelegramSearchPage( + val videos: List, + val answered: Boolean, + val hasMore: Boolean +) + +/** + * The phrase-by-phrase part of a Telegram lookup, apart from Android and TDLib so it can be tested. + * + * - Phrasings go out a pair at a time: global search is what Telegram rate-limits, and a burst of + * all of them made most go unanswered. + * - The core phrasings always run; the fallback ones only when the core found no matching file, + * and they stop at the first pair that does. + * - One unfiltered page per phrasing. When that page has more results behind it, newer text or + * photo posts may have pushed the videos off it (the search is newest-first), so the phrasing is + * repeated for videos and for documents — only then, to keep the request count down. + * - [Outcome.complete] is false as soon as any request went unanswered: a lookup with a gap is not + * an answer and must not be cached as one. + */ +class TelegramPhraseSearch( + private val fetchPage: suspend (query: String, filter: TelegramSearchFilter) -> TelegramSearchPage, + private val parallelQueries: Int = 2 +) { + data class Outcome( + /** Matching files, first-found order, one per (file name, size). */ + val matched: List, + val complete: Boolean, + val phrasingsSent: Int, + val requestsSent: Int, + val distinctFiles: Int + ) + + suspend fun run( + core: List, + fallback: List, + keep: (TelegramVideoMessage) -> Boolean, + matches: (TelegramVideoMessage) -> Boolean, + onMatchedGrew: suspend (List) -> Unit = {} + ): Outcome { + val seen = HashSet>() + val matched = LinkedHashMap, TelegramVideoMessage>() + var complete = true + var phrasingsSent = 0 + var requestsSent = 0 + + suspend fun searchBatch(batch: List) { + val results = coroutineScope { + batch.map { query -> async { searchPhrase(query) } }.awaitAll() + } + var grew = false + for (result in results) { + requestsSent += result.requests + if (!result.complete) complete = false + for (msg in result.videos) { + if (!keep(msg)) continue + val fileKey = msg.fileName to msg.fileSize + if (!seen.add(fileKey) || !matches(msg)) continue + matched[fileKey] = msg + grew = true + } + } + phrasingsSent += batch.size + if (grew) onMatchedGrew(matched.values.toList()) + } + + for (batch in core.chunked(parallelQueries)) searchBatch(batch) + if (matched.isEmpty()) { + for (batch in fallback.chunked(parallelQueries)) { + searchBatch(batch) + if (matched.isNotEmpty()) break + } + } + return Outcome(matched.values.toList(), complete, phrasingsSent, requestsSent, seen.size) + } + + private class PhraseResult(val videos: List, val complete: Boolean, val requests: Int) + + private suspend fun searchPhrase(query: String): PhraseResult { + val first = fetchPage(query, TelegramSearchFilter.ALL) + if (!first.answered) return PhraseResult(emptyList(), complete = false, requests = 1) + if (!first.hasMore) return PhraseResult(first.videos, complete = true, requests = 1) + // The unfiltered page was cut off: ask for the file types themselves. + val videos = first.videos.toMutableList() + var complete = true + for (filter in listOf(TelegramSearchFilter.VIDEO, TelegramSearchFilter.DOCUMENT)) { + val page = fetchPage(query, filter) + if (!page.answered) complete = false + videos += page.videos + } + return PhraseResult(videos, complete, requests = 3) + } +} + +/** + * Completed Telegram lookups by title/episode. Only a complete lookup is stored: a timed-out or + * interrupted one says nothing about what the groups hold, and caching it hid Telegram sources for + * hours (2h for a series). + */ +class TelegramResultCache(private val now: () -> Long = System::currentTimeMillis) { + private class Entry(val results: List, val expiresAt: Long) + + private val entries = java.util.concurrent.ConcurrentHashMap>() + + /** Returns whether [results] were stored. */ + fun put(key: String, results: List, complete: Boolean, ttlMs: Long): Boolean { + if (!complete) return false + entries[key] = Entry(results, now() + ttlMs) + return true + } + + fun get(key: String): List? = entries[key]?.takeIf { now() < it.expiresAt }?.results +} diff --git a/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramRepository.kt b/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramRepository.kt index 2bed6f617..4bf224453 100644 --- a/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramRepository.kt +++ b/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramRepository.kt @@ -191,74 +191,63 @@ class TelegramRepository @Inject constructor( } /** - * Searches globally across all chats (equivalent to Telethon's iter_messages(None, ...)) and - * keeps the video files. - * - * ONE request, with no type filter; videos and video documents are picked out below. It used to - * be two per phrase (Document, then Video), and global search is exactly what Telegram rate- - * limits: one episode lookup sent 26 of them, most of which TDLib then held back until the - * caller gave up (Special Ops S1E1, Sept 2026: 8 of 9 requests unanswered after 10s). + * One page of a global search across all chats (equivalent to Telethon's iter_messages(None, ...)), + * keeping the video files. [TelegramPhraseSearch] decides how many pages and which filters a + * phrasing needs; see there for why global search is kept to as few requests as possible. */ - suspend fun searchVideoMessages( + suspend fun searchPage( query: String, - limit: Int = 50 - ): List { - val filters = listOf(null) - val seen = mutableSetOf>() // dedupe by (fileName, fileSize) - val results = mutableListOf() - - for (filter in filters) { - // Global search answers in 1–9s (Special Ops S1E1, Sept 2026: "פרק 1" 7.9s, - // "lioness s01e01" 8.9s); the default 10s cut those off on a slower second try. - val result = client.sendRequest(TdApi.SearchMessages().also { req -> - req.chatList = null // null = search all chats (like Telethon's iter_messages(None)) - req.query = query - req.offset = "" - req.limit = limit - req.filter = filter - }, timeoutMs = SEARCH_REQUEST_TIMEOUT_MS) - val found = (result as? TdApi.FoundMessages) ?: continue + filter: TelegramSearchFilter, + limit: Int + ): TelegramSearchPage { + // Global search answers in 1–9s (Special Ops S1E1, Sept 2026: "פרק 1" 7.9s, + // "lioness s01e01" 8.9s); the default 10s cut those off on a slower second try. + val result = client.sendRequest(TdApi.SearchMessages().also { req -> + req.chatList = null // null = search all chats (like Telethon's iter_messages(None)) + req.query = query + req.offset = "" + req.limit = limit + req.filter = when (filter) { + TelegramSearchFilter.ALL -> null + TelegramSearchFilter.VIDEO -> TdApi.SearchMessagesFilterVideo() + TelegramSearchFilter.DOCUMENT -> TdApi.SearchMessagesFilterDocument() + } + }, timeoutMs = SEARCH_REQUEST_TIMEOUT_MS) + // No answer (timeout, or TDLib aborted it) is not an empty result — the caller must not + // treat the lookup as complete. + val found = result as? TdApi.FoundMessages + ?: return TelegramSearchPage(emptyList(), answered = false, hasMore = false) - for (msg in found.messages) { - when (val content = msg.content) { - is TdApi.MessageDocument -> { - val mime = content.document.mimeType - if (!mime.startsWith("video/") && mime != "application/x-matroska") continue - val key = content.document.fileName to content.document.document.size - if (seen.add(key)) { - results.add(TelegramVideoMessage( - messageId = msg.id, - chatId = msg.chatId, - fileName = content.document.fileName, - fileId = content.document.document.id, - fileSize = content.document.document.size, - duration = 0, - mimeType = mime, - caption = content.caption.text - )) - } - } - is TdApi.MessageVideo -> { - val key = content.video.fileName to content.video.video.size - if (seen.add(key)) { - results.add(TelegramVideoMessage( - messageId = msg.id, - chatId = msg.chatId, - fileName = content.video.fileName, - fileId = content.video.video.id, - fileSize = content.video.video.size, - duration = content.video.duration, - mimeType = content.video.mimeType, - caption = content.caption.text - )) - } - } - else -> continue + val videos = found.messages.mapNotNull { msg -> + when (val content = msg.content) { + is TdApi.MessageDocument -> { + val mime = content.document.mimeType + if (!mime.startsWith("video/") && mime != "application/x-matroska") return@mapNotNull null + TelegramVideoMessage( + messageId = msg.id, + chatId = msg.chatId, + fileName = content.document.fileName, + fileId = content.document.document.id, + fileSize = content.document.document.size, + duration = 0, + mimeType = mime, + caption = content.caption.text + ) } + is TdApi.MessageVideo -> TelegramVideoMessage( + messageId = msg.id, + chatId = msg.chatId, + fileName = content.video.fileName, + fileId = content.video.video.id, + fileSize = content.video.video.size, + duration = content.video.duration, + mimeType = content.video.mimeType, + caption = content.caption.text + ) + else -> null } - } - - return results + }.distinctBy { it.fileName to it.fileSize } + return TelegramSearchPage(videos, answered = true, hasMore = found.nextOffset.isNotEmpty()) } fun getStreamUrl(fileId: Int): String = proxy.getUrl(fileId) diff --git a/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramSourceResolver.kt b/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramSourceResolver.kt index 56ec06a97..8f9a8aa6e 100644 --- a/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramSourceResolver.kt +++ b/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramSourceResolver.kt @@ -16,7 +16,6 @@ import kotlinx.coroutines.Deferred import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.async -import kotlinx.coroutines.awaitAll import kotlinx.coroutines.coroutineScope import kotlinx.coroutines.currentCoroutineContext import kotlinx.coroutines.ensureActive @@ -60,8 +59,7 @@ class TelegramSourceResolver @Inject constructor( private const val CACHE_TTL_LONG_MS = 24 * 60 * 60 * 1_000L } - private data class CacheEntry(val results: List, val expiresAt: Long) - private val cache = ConcurrentHashMap() + private val cache = TelegramResultCache() /** * A running lookup: [found] holds every matching source so far (whole list, replaced as it @@ -94,10 +92,7 @@ class TelegramSourceResolver @Inject constructor( /** What an earlier completed lookup found for this title/episode, without searching. */ fun cachedResults(title: String, season: Int? = null, episode: Int? = null, imdbId: String = ""): List = - cache[cacheKey(imdbId, title, season, episode)] - ?.takeIf { System.currentTimeMillis() < it.expiresAt } - ?.results - .orEmpty() + cache.get(cacheKey(imdbId, title, season, episode)).orEmpty() /** * Playback is starting: stop every running search (a source list the user has already chosen @@ -153,9 +148,7 @@ class TelegramSourceResolver @Inject constructor( if (!repository.isAuthenticated()) return TelegramResolution(emptyList(), complete = true) val key = cacheKey(imdbId, title, season, episode) - cache[key]?.let { entry -> - if (System.currentTimeMillis() < entry.expiresAt) return TelegramResolution(entry.results, complete = true) - } + cache.get(key)?.let { return TelegramResolution(it, complete = true) } val search = inFlight.compute(key) { _, running -> running?.takeIf { it.result.isActive } ?: run { @@ -193,25 +186,22 @@ class TelegramSourceResolver @Inject constructor( found: MutableStateFlow> ): TelegramResolution = try { val startedAt = System.currentTimeMillis() - val finished = withTimeoutOrNull(SEARCH_TIMEOUT_MS) { + // null: the whole lookup ran out of time; false: it ended but a request went unanswered. + val complete = withTimeoutOrNull(SEARCH_TIMEOUT_MS) { resolveInternal(title, year, season, episode, imdbId, isMovie, found) - } != null + } ?: false val results = found.value Log.i( TAG, "Telegram lookup '$title'${season?.let { " S${it}E$episode" }.orEmpty()}: " + - "${results.size} source(s)" + (if (finished) "" else ", timed out") + + "${results.size} source(s)" + (if (complete) "" else ", incomplete (Telegram did not answer in time)") + " after ${System.currentTimeMillis() - startedAt}ms" ) - if (!finished) { - // Not cached: a timeout says Telegram was slow, not that the groups have nothing. - // Caching it hid the sources for hours (2h for a series). What was found is kept. - if (results.isEmpty()) notifyUser(context.getString(R.string.telegram_search_timed_out)) - TelegramResolution(results, complete = false) - } else { - cache[key] = CacheEntry(results, System.currentTimeMillis() + cacheTtl(year, isMovie)) - TelegramResolution(results, complete = true) - } + // Not cached when incomplete: a timeout says Telegram was slow, not that the groups have + // nothing — caching it hid the sources for hours (2h for a series). What was found is kept. + cache.put(key, results, complete, cacheTtl(year, isMovie)) + if (!complete && results.isEmpty()) notifyUser(context.getString(R.string.telegram_search_timed_out)) + TelegramResolution(results, complete) } catch (e: TelegramApiException) { Log.w(TAG, "Telegram API error for '$title': ${e.message}") if (found.value.isEmpty()) notifyUser(friendlyError(e.message)) @@ -249,7 +239,7 @@ class TelegramSourceResolver @Inject constructor( imdbId: String, isMovie: Boolean, found: MutableStateFlow> - ) { + ): Boolean { val excludedIds = repository.getExcludedChatIds().first() // Read content language from SharedPreferences (same store SettingsViewModel writes to). // Avoids a DI cycle: StreamRepository → TelegramSourceResolver → MediaRepository → StreamRepository. @@ -269,70 +259,46 @@ class TelegramSourceResolver @Inject constructor( matcher.buildMovieQueries(title, year, localizedTitle, englishTitle, originalTitle) val (core, fallback) = splitCoreQueries(queries, season, episode) - val seen = mutableSetOf>() - // Matching sources by file, built once: the stream URL registers the file with the proxy. - val matched = LinkedHashMap, StreamSource>() + // Stream sources by file, built once: the stream URL registers the file with the proxy. + val sources = HashMap, StreamSource>() val searchStartedAt = System.currentTimeMillis() - fun matchesRequest(msg: TelegramVideoMessage): Boolean = matcher.score( - fileName = msg.fileName, - caption = msg.caption, - title = title, - localizedTitle = localizedTitle, - englishTitle = englishTitle, - originalTitle = originalTitle, - year = year, - season = season, - episode = episode - ) >= SCORE_THRESHOLD - - suspend fun searchBatch(batch: List) { - val messages = coroutineScope { - batch.map { query -> - async { - try { - repository.searchVideoMessages(query, MAX_RESULTS) - .filter { it.chatId !in excludedIds } - } catch (e: TelegramApiException) { - throw e - } catch (e: kotlinx.coroutines.CancellationException) { - throw e - } catch (e: Exception) { - Log.e(TAG, "Search failed for '$query'", e) - emptyList() - } - } - }.awaitAll().flatten() - } - var grew = false - for (msg in messages) { - val fileKey = msg.fileName to msg.fileSize - if (!seen.add(fileKey) || !matchesRequest(msg)) continue - matched[fileKey] = toStreamSource(msg) - grew = true - } + val outcome = TelegramPhraseSearch( + fetchPage = { query, filter -> repository.searchPage(query, filter, MAX_RESULTS) }, + parallelQueries = MAX_PARALLEL_QUERIES + ).run( + core = core, + fallback = fallback, + keep = { it.chatId !in excludedIds }, + matches = { msg -> + matcher.score( + fileName = msg.fileName, + caption = msg.caption, + title = title, + localizedTitle = localizedTitle, + englishTitle = englishTitle, + originalTitle = originalTitle, + year = year, + season = season, + episode = episode + ) >= SCORE_THRESHOLD + }, // Shown as soon as they are found; the lookup keeps going. - if (grew) found.value = sortSources(matched.values, langCode) - } - - var sent = 0 - for (batch in core.chunked(MAX_PARALLEL_QUERIES)) { - searchBatch(batch) - sent += batch.size - } - if (matched.isEmpty()) { - for (batch in fallback.chunked(MAX_PARALLEL_QUERIES)) { - searchBatch(batch) - sent += batch.size - if (matched.isNotEmpty()) break + onMatchedGrew = { matched -> + found.value = sortSources( + matched.map { msg -> sources.getOrPut(msg.fileName to msg.fileSize) { toStreamSource(msg) } }, + langCode + ) } - } + ) Log.i( TAG, - "Telegram search: $sent of ${queries.size} phrasings sent (${core.size} core, " + - "$MAX_PARALLEL_QUERIES at a time) in ${System.currentTimeMillis() - searchStartedAt}ms, " + - "${seen.size} distinct files, ${matched.size} matching" + "Telegram search: ${outcome.phrasingsSent} of ${queries.size} phrasings sent (${core.size} core, " + + "$MAX_PARALLEL_QUERIES at a time, ${outcome.requestsSent} requests) in " + + "${System.currentTimeMillis() - searchStartedAt}ms, ${outcome.distinctFiles} distinct files, " + + "${outcome.matched.size} matching" + if (outcome.complete) "" else ", some requests unanswered" ) + return outcome.complete } private fun toStreamSource(msg: TelegramVideoMessage): StreamSource { diff --git a/app/src/test/kotlin/com/arflix/tv/data/repository/TelegramCachedSourcesTest.kt b/app/src/test/kotlin/com/arflix/tv/data/repository/TelegramCachedSourcesTest.kt new file mode 100644 index 000000000..5adf2bc62 --- /dev/null +++ b/app/src/test/kotlin/com/arflix/tv/data/repository/TelegramCachedSourcesTest.kt @@ -0,0 +1,50 @@ +package com.arflix.tv.data.repository + +import com.arflix.tv.data.model.StreamSource +import org.junit.Assert.assertEquals +import org.junit.Assert.assertSame +import org.junit.Test + +/** + * Reopening Sources after a click-search: the saved source list holds only addon results, the + * Telegram files live in the Telegram lookup's cache, and every saved-list return path merges them + * back (resolveMovieStreamsProgressive / resolveEpisodeStreamsProgressive). + */ +class TelegramCachedSourcesTest { + + private fun addon(name: String) = StreamSource( + source = name, addonName = "Torrentio", addonId = "torrentio", quality = "1080p", size = "2 GB", + url = "https://example.com/$name" + ) + + private fun telegram(name: String) = StreamSource( + source = name, addonName = "Telegram", addonId = "telegram_native", quality = "1080p", size = "2 GB", + url = "http://localhost:1234/file/${name.hashCode()}" + ) + + @Test + fun `a saved addon-only list gets the click-search results back`() { + val saved = listOf(addon("a.mkv"), addon("b.mkv")) + val found = listOf(telegram("Show.S01E04.mkv")) + assertEquals(saved + found, withCachedTelegramSources(saved, found)) + } + + @Test + fun `with Telegram as the only provider the list is the click-search results`() { + val found = listOf(telegram("Show.S01E04.mkv"), telegram("Show.S01E04.720p.mkv")) + assertEquals(found, withCachedTelegramSources(emptyList(), found)) + } + + @Test + fun `a list that already has the Telegram files does not list them twice`() { + val found = listOf(telegram("Show.S01E04.mkv")) + val saved = listOf(addon("a.mkv")) + found + assertEquals(saved, withCachedTelegramSources(saved, found)) + } + + @Test + fun `nothing cached leaves the saved list untouched`() { + val saved = listOf(addon("a.mkv")) + assertSame(saved, withCachedTelegramSources(saved, emptyList())) + } +} diff --git a/app/src/test/kotlin/com/arflix/tv/data/telegram/TelegramPhraseSearchTest.kt b/app/src/test/kotlin/com/arflix/tv/data/telegram/TelegramPhraseSearchTest.kt new file mode 100644 index 000000000..4725335de --- /dev/null +++ b/app/src/test/kotlin/com/arflix/tv/data/telegram/TelegramPhraseSearchTest.kt @@ -0,0 +1,185 @@ +package com.arflix.tv.data.telegram + +import kotlinx.coroutines.runBlocking +import org.junit.Assert.assertEquals +import org.junit.Assert.assertFalse +import org.junit.Assert.assertTrue +import org.junit.Test + +class TelegramPhraseSearchTest { + + private fun video(name: String, chatId: Long = 1L, size: Long = name.length * 1_000L) = TelegramVideoMessage( + messageId = size, + chatId = chatId, + fileName = name, + fileId = name.hashCode(), + fileSize = size, + duration = 0, + mimeType = "video/x-matroska", + caption = "" + ) + + private val episodeMatcher: (TelegramVideoMessage) -> Boolean = { it.fileName.contains("S01E04") } + + /** Answers from [pages] by (query, filter), recording every request. */ + private class FakeTelegram(private val pages: Map, TelegramSearchPage>) { + val requests = mutableListOf>() + suspend fun fetch(query: String, filter: TelegramSearchFilter): TelegramSearchPage { + synchronized(requests) { requests += query to filter } + return pages[query to filter] ?: TelegramSearchPage(emptyList(), answered = true, hasMore = false) + } + } + + private fun answered(vararg videos: TelegramVideoMessage, hasMore: Boolean = false) = + TelegramSearchPage(videos.toList(), answered = true, hasMore = hasMore) + + private val unanswered = TelegramSearchPage(emptyList(), answered = false, hasMore = false) + + @Test + fun `a request Telegram did not answer makes the lookup incomplete but keeps what was found`() = runBlocking { + val telegram = FakeTelegram( + mapOf( + ("show ע1פ4" to TelegramSearchFilter.ALL) to unanswered, + ("show ע1 פ4" to TelegramSearchFilter.ALL) to answered(video("Show.S01E04.1080p.mkv")) + ) + ) + val outcome = TelegramPhraseSearch(telegram::fetch).run( + core = listOf("show ע1פ4", "show ע1 פ4"), + fallback = emptyList(), + keep = { true }, + matches = episodeMatcher + ) + assertFalse("a lookup with an unanswered request must not be cached as an answer", outcome.complete) + assertEquals(listOf("Show.S01E04.1080p.mkv"), outcome.matched.map { it.fileName }) + } + + @Test + fun `a lookup where every request was answered is complete even with no matches`() = runBlocking { + val telegram = FakeTelegram(emptyMap()) + val outcome = TelegramPhraseSearch(telegram::fetch).run( + core = listOf("show ע1פ4", "show ע1 פ4"), + fallback = emptyList(), + keep = { true }, + matches = episodeMatcher + ) + assertTrue(outcome.complete) + assertTrue(outcome.matched.isEmpty()) + } + + @Test + fun `a first page cut off by non-video posts is repeated for videos and documents`() = runBlocking { + // The unfiltered page holds only text/photo posts (no videos) and has more behind it — + // the matching file sits further back and only the typed searches reach it. + val telegram = FakeTelegram( + mapOf( + ("show פרק 4" to TelegramSearchFilter.ALL) to answered(hasMore = true), + ("show פרק 4" to TelegramSearchFilter.VIDEO) to answered(video("Show.S01E04.720p.mp4")), + ("show פרק 4" to TelegramSearchFilter.DOCUMENT) to answered(video("Show.S01E04.1080p.mkv")) + ) + ) + val outcome = TelegramPhraseSearch(telegram::fetch).run( + core = listOf("show פרק 4"), + fallback = emptyList(), + keep = { true }, + matches = episodeMatcher + ) + assertEquals( + listOf( + "show פרק 4" to TelegramSearchFilter.ALL, + "show פרק 4" to TelegramSearchFilter.VIDEO, + "show פרק 4" to TelegramSearchFilter.DOCUMENT + ), + telegram.requests + ) + assertEquals(setOf("Show.S01E04.720p.mp4", "Show.S01E04.1080p.mkv"), outcome.matched.map { it.fileName }.toSet()) + assertTrue(outcome.complete) + } + + @Test + fun `a page with nothing behind it costs one request`() = runBlocking { + val telegram = FakeTelegram( + mapOf(("show ע1פ4" to TelegramSearchFilter.ALL) to answered(video("Show.S01E04.mkv"))) + ) + TelegramPhraseSearch(telegram::fetch).run( + core = listOf("show ע1פ4"), + fallback = emptyList(), + keep = { true }, + matches = episodeMatcher + ) + assertEquals(listOf("show ע1פ4" to TelegramSearchFilter.ALL), telegram.requests) + } + + @Test + fun `an unanswered first page is not retried with filters`() = runBlocking { + // Telegram is throttling: more requests only queue behind the held ones. + val telegram = FakeTelegram(mapOf(("show ע1פ4" to TelegramSearchFilter.ALL) to unanswered)) + val outcome = TelegramPhraseSearch(telegram::fetch).run( + core = listOf("show ע1פ4"), + fallback = emptyList(), + keep = { true }, + matches = episodeMatcher + ) + assertEquals(1, telegram.requests.size) + assertFalse(outcome.complete) + } + + @Test + fun `fallback phrasings run only when the core found nothing and stop at the first match`() = runBlocking { + val telegram = FakeTelegram( + mapOf(("show s1 e4" to TelegramSearchFilter.ALL) to answered(video("Show S01E04.mkv"))) + ) + val outcome = TelegramPhraseSearch(telegram::fetch, parallelQueries = 2).run( + core = listOf("show ע1פ4", "show ע1 פ4"), + fallback = listOf("show s1e4", "show s1 e4", "show s01 e04", "show פ4"), + keep = { true }, + matches = episodeMatcher + ) + assertEquals(listOf("Show S01E04.mkv"), outcome.matched.map { it.fileName }) + // Two core + the first fallback pair; the second fallback pair is never sent. + assertEquals(4, outcome.phrasingsSent) + assertFalse(telegram.requests.any { it.first == "show s01 e04" || it.first == "show פ4" }) + + val coreHit = FakeTelegram( + mapOf(("show ע1פ4" to TelegramSearchFilter.ALL) to answered(video("Show.S01E04.mkv"))) + ) + TelegramPhraseSearch(coreHit::fetch).run( + core = listOf("show ע1פ4"), + fallback = listOf("show s1e4"), + keep = { true }, + matches = episodeMatcher + ) + assertFalse(coreHit.requests.any { it.first == "show s1e4" }) + } + + @Test + fun `excluded chats and non-matching files are dropped`() = runBlocking { + val telegram = FakeTelegram( + mapOf( + ("show ע1פ4" to TelegramSearchFilter.ALL) to answered( + video("Show.S01E04.excluded.mkv", chatId = 99L), + video("Show.S01E05.mkv"), + video("Show.S01E04.kept.mkv") + ) + ) + ) + val outcome = TelegramPhraseSearch(telegram::fetch).run( + core = listOf("show ע1פ4"), + fallback = emptyList(), + keep = { it.chatId != 99L }, + matches = episodeMatcher + ) + assertEquals(listOf("Show.S01E04.kept.mkv"), outcome.matched.map { it.fileName }) + } + + @Test + fun `result cache keeps only complete lookups and expires them`() { + var now = 1_000L + val cache = TelegramResultCache(now = { now }) + assertFalse(cache.put("tt1:1:4", listOf("partial"), complete = false, ttlMs = 60_000L)) + assertEquals(null, cache.get("tt1:1:4")) + assertTrue(cache.put("tt1:1:4", listOf("a", "b"), complete = true, ttlMs = 60_000L)) + assertEquals(listOf("a", "b"), cache.get("tt1:1:4")) + now += 60_000L + assertEquals(null, cache.get("tt1:1:4")) + } +} From c3a9098a28b1c495dc9872a3cd95d8bb9385ef12 Mon Sep 17 00:00:00 2001 From: Arvin Date: Sun, 27 Sep 2026 15:15:56 +0200 Subject: [PATCH 3/3] fix(telegram): retain results while sibling searches are pending --- .../tv/data/telegram/TelegramPhraseSearch.kt | 44 ++++++++++------- .../data/telegram/TelegramPhraseSearchTest.kt | 47 +++++++++++++++++++ 2 files changed, 74 insertions(+), 17 deletions(-) diff --git a/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramPhraseSearch.kt b/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramPhraseSearch.kt index e5d719a4e..d63fae2f0 100644 --- a/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramPhraseSearch.kt +++ b/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramPhraseSearch.kt @@ -3,6 +3,8 @@ package com.arflix.tv.data.telegram import kotlinx.coroutines.async import kotlinx.coroutines.awaitAll import kotlinx.coroutines.coroutineScope +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock /** Which messages a global search page asks Telegram for. */ enum class TelegramSearchFilter { ALL, VIDEO, DOCUMENT } @@ -58,25 +60,29 @@ class TelegramPhraseSearch( var complete = true var phrasingsSent = 0 var requestsSent = 0 + val resultsMutex = Mutex() + + suspend fun publish(videos: List) = resultsMutex.withLock { + var grew = false + for (msg in videos) { + if (!keep(msg)) continue + val fileKey = msg.fileName to msg.fileSize + if (!seen.add(fileKey) || !matches(msg)) continue + matched[fileKey] = msg + grew = true + } + if (grew) onMatchedGrew(matched.values.toList()) + } suspend fun searchBatch(batch: List) { val results = coroutineScope { - batch.map { query -> async { searchPhrase(query) } }.awaitAll() + batch.map { query -> async { searchPhrase(query, ::publish) } }.awaitAll() } - var grew = false for (result in results) { requestsSent += result.requests if (!result.complete) complete = false - for (msg in result.videos) { - if (!keep(msg)) continue - val fileKey = msg.fileName to msg.fileSize - if (!seen.add(fileKey) || !matches(msg)) continue - matched[fileKey] = msg - grew = true - } } phrasingsSent += batch.size - if (grew) onMatchedGrew(matched.values.toList()) } for (batch in core.chunked(parallelQueries)) searchBatch(batch) @@ -89,21 +95,25 @@ class TelegramPhraseSearch( return Outcome(matched.values.toList(), complete, phrasingsSent, requestsSent, seen.size) } - private class PhraseResult(val videos: List, val complete: Boolean, val requests: Int) + private class PhraseResult(val complete: Boolean, val requests: Int) - private suspend fun searchPhrase(query: String): PhraseResult { + private suspend fun searchPhrase( + query: String, + publish: suspend (List) -> Unit + ): PhraseResult { val first = fetchPage(query, TelegramSearchFilter.ALL) - if (!first.answered) return PhraseResult(emptyList(), complete = false, requests = 1) - if (!first.hasMore) return PhraseResult(first.videos, complete = true, requests = 1) + if (!first.answered) return PhraseResult(complete = false, requests = 1) + // Preserve each answered page even if a sibling query or a filtered follow-up stalls. + publish(first.videos) + if (!first.hasMore) return PhraseResult(complete = true, requests = 1) // The unfiltered page was cut off: ask for the file types themselves. - val videos = first.videos.toMutableList() var complete = true for (filter in listOf(TelegramSearchFilter.VIDEO, TelegramSearchFilter.DOCUMENT)) { val page = fetchPage(query, filter) if (!page.answered) complete = false - videos += page.videos + publish(page.videos) } - return PhraseResult(videos, complete, requests = 3) + return PhraseResult(complete, requests = 3) } } diff --git a/app/src/test/kotlin/com/arflix/tv/data/telegram/TelegramPhraseSearchTest.kt b/app/src/test/kotlin/com/arflix/tv/data/telegram/TelegramPhraseSearchTest.kt index 4725335de..6b58945e6 100644 --- a/app/src/test/kotlin/com/arflix/tv/data/telegram/TelegramPhraseSearchTest.kt +++ b/app/src/test/kotlin/com/arflix/tv/data/telegram/TelegramPhraseSearchTest.kt @@ -1,6 +1,9 @@ package com.arflix.tv.data.telegram import kotlinx.coroutines.runBlocking +import kotlinx.coroutines.awaitCancellation +import kotlinx.coroutines.test.runTest +import kotlinx.coroutines.withTimeoutOrNull import org.junit.Assert.assertEquals import org.junit.Assert.assertFalse import org.junit.Assert.assertTrue @@ -35,6 +38,50 @@ class TelegramPhraseSearchTest { private val unanswered = TelegramSearchPage(emptyList(), answered = false, hasMore = false) + @Test + fun `completed query publishes results even when its sibling times out`() = runTest { + val found = mutableListOf>() + var stalledCancelled = false + val outcome = withTimeoutOrNull(100) { + TelegramPhraseSearch(fetchPage = { query, _ -> + if (query == "slow") { + try { awaitCancellation() } finally { stalledCancelled = true } + } else answered(video("Show.S01E04.mkv")) + }).run( + core = listOf("slow", "fast"), fallback = emptyList(), + keep = { true }, matches = episodeMatcher, + onMatchedGrew = { found += it } + ) + } + assertEquals(null, outcome) + assertTrue(stalledCancelled) + assertEquals(listOf("Show.S01E04.mkv"), found.last().map { it.fileName }) + } + + @Test + fun `answered pages survive a stalled filtered follow-up without duplicate results`() = runTest { + val found = mutableListOf>() + val requests = mutableListOf() + val outcome = withTimeoutOrNull(100) { + TelegramPhraseSearch(fetchPage = { _, filter -> + requests += filter + when (filter) { + TelegramSearchFilter.ALL -> answered(video("Show.S01E04.mkv"), hasMore = true) + TelegramSearchFilter.VIDEO -> answered(video("Show.S01E04.mkv"), video("Show.S01E04.720p.mp4")) + TelegramSearchFilter.DOCUMENT -> awaitCancellation() + } + }).run( + core = listOf("show"), fallback = emptyList(), + keep = { true }, matches = episodeMatcher, + onMatchedGrew = { found += it } + ) + } + assertEquals(null, outcome) + assertEquals(listOf(TelegramSearchFilter.ALL, TelegramSearchFilter.VIDEO, TelegramSearchFilter.DOCUMENT), requests) + assertEquals(2, found.size) + assertEquals(listOf("Show.S01E04.mkv", "Show.S01E04.720p.mp4"), found.last().map { it.fileName }) + } + @Test fun `a request Telegram did not answer makes the lookup incomplete but keeps what was found`() = runBlocking { val telegram = FakeTelegram(