Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down Expand Up @@ -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)))
Expand Down Expand Up @@ -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")
Expand Down
241 changes: 195 additions & 46 deletions app/src/main/kotlin/com/arflix/tv/data/repository/StreamRepository.kt

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -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) },
Expand Down
Original file line number Diff line number Diff line change
@@ -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<TelegramVideoMessage>,
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<TelegramVideoMessage>,
val complete: Boolean,
val phrasingsSent: Int,
val requestsSent: Int,
val distinctFiles: Int
)

suspend fun run(
core: List<String>,
fallback: List<String>,
keep: (TelegramVideoMessage) -> Boolean,
matches: (TelegramVideoMessage) -> Boolean,
onMatchedGrew: suspend (List<TelegramVideoMessage>) -> Unit = {}
): Outcome {
val seen = HashSet<Pair<String, Long>>()
val matched = LinkedHashMap<Pair<String, Long>, TelegramVideoMessage>()
var complete = true
var phrasingsSent = 0
var requestsSent = 0
val resultsMutex = Mutex()

suspend fun publish(videos: List<TelegramVideoMessage>) = 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<String>) {
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<TelegramVideoMessage>) -> 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<T>(private val now: () -> Long = System::currentTimeMillis) {
private class Entry<T>(val results: List<T>, val expiresAt: Long)

private val entries = java.util.concurrent.ConcurrentHashMap<String, Entry<T>>()

/** Returns whether [results] were stored. */
fun put(key: String, results: List<T>, complete: Boolean, ttlMs: Long): Boolean {
if (!complete) return false
entries[key] = Entry(results, now() + ttlMs)
return true
}

fun get(key: String): List<T>? = entries[key]?.takeIf { now() < it.expiresAt }?.results
}
139 changes: 79 additions & 60 deletions app/src/main/kotlin/com/arflix/tv/data/telegram/TelegramRepository.kt
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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")

Expand Down Expand Up @@ -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<TelegramVideoMessage> {
val filters = listOf(
TdApi.SearchMessagesFilterDocument(),
TdApi.SearchMessagesFilterVideo()
)
val seen = mutableSetOf<Pair<String, Long>>() // dedupe by (fileName, fileSize)
val results = mutableListOf<TelegramVideoMessage>()

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<Boolean> =
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<Set<Long>> =
context.telegramDataStore.data.map { prefs ->
prefs[KEY_EXCLUDED_CHATS]
Expand Down
Loading
Loading