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 eda6ff0e5..ebe1fdd22 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))) @@ -1749,6 +1753,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 1e7bfe040..09afbff53 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(), @@ -1759,6 +1771,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 @@ -2334,16 +2356,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 + } } } @@ -2359,11 +2387,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) { @@ -2394,13 +2424,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 } @@ -2418,7 +2458,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, @@ -2433,7 +2473,6 @@ class StreamRepository @Inject constructor( } val prioritizedAddons = prioritizeStreamingAddons(streamAddons) - val telegramEnabled = telegramSourceResolver.isEnabled() && streamIntegrationRepository.isIntegrationEnabled(StreamIntegrationType.TELEGRAM) if (prioritizedAddons.isEmpty() && !telegramEnabled) { Log.w( TAG, @@ -2447,19 +2486,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 } @@ -2470,8 +2509,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() { @@ -2485,14 +2528,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", @@ -2546,12 +2591,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 -> @@ -2601,15 +2649,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() @@ -2678,6 +2743,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, @@ -2690,6 +2793,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) { @@ -2930,11 +3042,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 + } } } } @@ -2963,7 +3079,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 } @@ -3006,13 +3122,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 } @@ -3022,7 +3148,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, @@ -3033,7 +3159,6 @@ class StreamRepository @Inject constructor( } val prioritizedAddons = prioritizeStreamingAddons(streamAddons) - val telegramEnabled = telegramSourceResolver.isEnabled() && streamIntegrationRepository.isIntegrationEnabled(StreamIntegrationType.TELEGRAM) if (prioritizedAddons.isEmpty() && !telegramEnabled) { Log.w( TAG, @@ -3047,12 +3172,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 } @@ -3063,8 +3188,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() { @@ -3078,8 +3207,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( @@ -3146,8 +3277,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, @@ -3155,10 +3286,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 -> @@ -3222,22 +3356,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/TelegramPhraseSearch.kt b/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramPhraseSearch.kt new file mode 100644 index 000000000..d63fae2f0 --- /dev/null +++ b/app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramPhraseSearch.kt @@ -0,0 +1,138 @@ +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 } + +/** + * 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 + 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, ::publish) } }.awaitAll() + } + for (result in results) { + requestsSent += result.requests + if (!result.complete) complete = false + } + phrasingsSent += batch.size + } + + 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 complete: Boolean, val requests: Int) + + private suspend fun searchPhrase( + query: String, + publish: suspend (List) -> Unit + ): PhraseResult { + val first = fetchPage(query, TelegramSearchFilter.ALL) + 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. + var complete = true + for (filter in listOf(TelegramSearchFilter.VIDEO, TelegramSearchFilter.DOCUMENT)) { + val page = fetchPage(query, filter) + if (!page.answered) complete = false + publish(page.videos) + } + return PhraseResult(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 39bf9a8ed..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 @@ -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,74 +191,82 @@ 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. + * 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( - TdApi.SearchMessagesFilterDocument(), - TdApi.SearchMessagesFilterVideo() - ) - val seen = mutableSetOf>() // dedupe by (fileName, fileSize) - val results = mutableListOf() - - for (filter in filters) { - 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 - }) - 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) + /** + * "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..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 @@ -10,15 +10,31 @@ 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,14 +45,41 @@ 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 } - 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 + * 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 ?: ""}" @@ -44,6 +87,34 @@ class TelegramSourceResolver @Inject constructor( 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.get(cacheKey(imdbId, title, season, episode)).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,39 +130,116 @@ 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 - } + cache.get(key)?.let { return TelegramResolution(it, 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() } - - 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() } } + 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() + // 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) + } ?: false + val results = found.value + Log.i( + TAG, + "Telegram lookup '$title'${season?.let { " S${it}E$episode" }.orEmpty()}: " + + "${results.size} source(s)" + (if (complete) "" else ", incomplete (Telegram did not answer in time)") + + " after ${System.currentTimeMillis() - startedAt}ms" + ) + // 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)) + 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( title: String, year: Int?, season: Int?, episode: Int?, imdbId: String, - isMovie: Boolean - ): List { + 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. @@ -109,35 +257,21 @@ 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() - } - } - }.awaitAll().flatten().forEach { msg -> - if (seen.add(msg.fileName to msg.fileSize)) allMessages.add(msg) - } - } + // Stream sources by file, built once: the stream URL registers the file with the proxy. + val sources = HashMap, StreamSource>() + val searchStartedAt = System.currentTimeMillis() - return allMessages - .mapNotNull { msg -> - val score = matcher.score( + 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, @@ -147,42 +281,60 @@ class TelegramSourceResolver @Inject constructor( 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() } + ) >= SCORE_THRESHOLD + }, + // Shown as soon as they are found; the lookup keeps going. + onMatchedGrew = { matched -> + found.value = sortSources( + matched.map { msg -> sources.getOrPut(msg.fileName to msg.fileSize) { toStreamSource(msg) } }, + langCode ) } - .sortedWith( - compareByDescending { langCode == "he" && matcher.isHebrew(it.source) } - .thenByDescending { qualityTier(it.quality) } - .thenByDescending { it.sizeBytes ?: 0L } - ) + ) + Log.i( + TAG, + "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 { + 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 +344,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 1aa609bd8..8b55bf4ff 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 @@ -104,6 +104,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, @@ -213,7 +215,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 && @@ -1967,6 +1980,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() @@ -1986,8 +2051,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 @@ -2052,6 +2119,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 5cd9fae8d..4ec4e9f89 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 @@ -654,6 +654,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, @@ -673,6 +684,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 dc4b27892..16f760cac 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? 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..6b58945e6 --- /dev/null +++ b/app/src/test/kotlin/com/arflix/tv/data/telegram/TelegramPhraseSearchTest.kt @@ -0,0 +1,232 @@ +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 +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 `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( + 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")) + } +}