cleanup logs

This commit is contained in:
tapframe 2026-06-04 16:25:51 +05:30
parent 20b1fd5408
commit 13fe6090c6
8 changed files with 80 additions and 337 deletions

View file

@ -46,6 +46,7 @@ private const val STREAMING_HANDSHAKE_TIMEOUT_MS = 3_000
private const val STREAMING_DISCONNECT_TIMEOUT_SECONDS = 120
private const val STREAMING_READ_AHEAD_PERCENT = 95
private const val STREAMING_PRELOAD_CACHE_PERCENT = 50
private const val STREAM_SCREEN_WARMUP_COOLDOWN_MS = 10_000L
private val VIDEO_EXTENSIONS = setOf("mkv", "mp4", "avi", "webm", "ts", "m4v", "mov", "wmv", "flv")
actual object P2pStreamingEngine {
@ -57,6 +58,7 @@ actual object P2pStreamingEngine {
private var statsJob: Job? = null
private var preloadJob: Job? = null
private var warmupJob: Job? = null
private var warmupCooldownJob: Job? = null
private var currentHash: String? = null
private var streamGeneration = 0L
private var appContext: Context? = null
@ -67,73 +69,86 @@ actual object P2pStreamingEngine {
fun initialize(context: Context) {
appContext = context.applicationContext
binary.initialize(context.applicationContext)
if (P2pSettingsRepository.isVisible && P2pSettingsStorage.loadP2pEnabled() == true) {
warmup()
}
}
actual fun warmup() {
if (appContext == null) return
synchronized(lifecycleLock) {
warmupCooldownJob?.cancel()
warmupCooldownJob = null
if (warmupJob?.isActive == true) return
warmupJob = scope.launch {
val startedAt = System.currentTimeMillis()
Log.d(TAG, "warmup: starting TorrServer")
try {
binary.start()
if (api.ensureStreamingSettings().changed) {
Log.d(TAG, "warmup: restarting TorrServer to apply streaming settings")
binary.stop()
binary.start()
api.ensureStreamingSettings()
}
Log.d(TAG, "warmup: TorrServer ready in ${System.currentTimeMillis() - startedAt}ms")
} catch (e: CancellationException) {
throw e
} catch (e: Exception) {
Log.w(TAG, "warmup: TorrServer failed after ${System.currentTimeMillis() - startedAt}ms", e)
Log.w(TAG, "TorrServer warmup failed", e)
}
}
}
}
actual fun cooldownWarmup() {
synchronized(lifecycleLock) {
if (currentHash != null) {
return
}
warmupCooldownJob?.cancel()
warmupCooldownJob = scope.launch {
delay(STREAM_SCREEN_WARMUP_COOLDOWN_MS)
val warmup = synchronized(lifecycleLock) {
if (currentHash != null) {
warmupCooldownJob = null
return@launch
}
val job = warmupJob
warmupJob = null
warmupCooldownJob = null
job
}
try {
warmup?.join()
} catch (_: CancellationException) {
}
val shouldStop = synchronized(lifecycleLock) { currentHash == null }
if (!shouldStop) {
return@launch
}
try {
binary.stop()
} catch (e: Exception) {
Log.w(TAG, "Failed to stop idle TorrServer", e)
}
}
}
}
actual suspend fun startStream(request: P2pStreamRequest): String = withContext(Dispatchers.IO) {
val startAt = System.currentTimeMillis()
val requestedHash = request.infoHash.trim().takeIf { it.isNotEmpty() }
?: throw P2pStreamingException("Missing torrent info hash")
val detached = beginStreamGeneration()
val generation = detached.generation
detached.hash?.let(::scheduleIdleDrop)
_state.value = P2pStreamingState.Connecting
Log.d(
TAG,
"startStream[$generation]: request hash=$requestedHash fileIdx=${request.fileIdx} " +
"filename=${request.filename.orEmpty()} magnet=${!request.magnetUri.isNullOrBlank()} " +
"trackers=${request.trackers.size} detachedHash=${detached.hash}",
)
var attachedHash: String? = null
try {
val warmupWaitAt = System.currentTimeMillis()
cancelWarmupCooldown()
awaitWarmup()
Log.d(TAG, "startStream[$generation]: warmup wait ${System.currentTimeMillis() - warmupWaitAt}ms")
val binaryStartAt = System.currentTimeMillis()
binary.start()
Log.d(TAG, "startStream[$generation]: binary ready ${System.currentTimeMillis() - binaryStartAt}ms")
ensureCurrentGeneration(generation)
val settingsAt = System.currentTimeMillis()
val settingsResult = api.ensureStreamingSettings()
if (settingsResult.changed) {
Log.d(TAG, "startStream[$generation]: restarting TorrServer to apply streaming settings")
binary.stop()
binary.start()
api.ensureStreamingSettings()
}
Log.d(
TAG,
"startStream[$generation]: streaming settings ready ${System.currentTimeMillis() - settingsAt}ms " +
"changed=${settingsResult.changed}",
)
ensureCurrentGeneration(generation)
val magnetLink = buildMagnetUri(
@ -141,12 +156,9 @@ actual object P2pStreamingEngine {
magnetUri = request.magnetUri,
extraTrackers = request.trackers,
)
Log.d(TAG, "startStream[$generation]: magnet ${summarizeMagnet(magnetLink)}")
val addAt = System.currentTimeMillis()
val hash = api.addTorrent(magnetLink)
?: throw P2pStreamingException("Failed to add torrent")
Log.d(TAG, "startStream[$generation]: addTorrent hash=$hash in ${System.currentTimeMillis() - addAt}ms")
attachedHash = hash
cancelIdleDrop(hash)
if (!attachTorrentIfCurrent(generation, hash)) {
@ -157,21 +169,13 @@ actual object P2pStreamingEngine {
val requestedName = request.filename?.trim()?.takeIf { it.isNotEmpty() }
val useEngineFileSelector = requestedName != null || request.fileIdx != null
val resolvedIdx = if (useEngineFileSelector) {
Log.d(
TAG,
"startStream[$generation]: using TorrServer selector filename=${requestedName.orEmpty()} " +
"fileIdx=${request.fileIdx ?: -1}",
)
null
} else {
val resolveAt = System.currentTimeMillis()
resolveFileIndex(
hash = hash,
requestedIdx = request.fileIdx,
filename = request.filename,
).also {
Log.d(TAG, "startStream[$generation]: resolved file index=$it in ${System.currentTimeMillis() - resolveAt}ms")
}
)
}
ensureCurrentGeneration(generation)
@ -184,7 +188,6 @@ actual object P2pStreamingEngine {
magnetLink = magnetLink,
selector = streamSelector,
)
Log.d(TAG, "startStream[$generation]: streamUrl=${summarizeStreamUrl(streamUrl)}")
startPreload(
hash = hash,
@ -195,11 +198,6 @@ actual object P2pStreamingEngine {
startStatsPolling(
hash = hash,
generation = generation,
fileSelection = ActiveFileSelection(
requestedIdx = request.fileIdx,
requestedFilename = request.filename,
resolvedIdx = resolvedIdx,
),
)
ensureCurrentGeneration(generation)
@ -213,14 +211,12 @@ actual object P2pStreamingEngine {
totalProgress = 0f,
)
Log.d(TAG, "startStream[$generation]: ready in ${System.currentTimeMillis() - startAt}ms")
streamUrl
} catch (e: CancellationException) {
Log.d(TAG, "startStream[$generation]: cancelled after ${System.currentTimeMillis() - startAt}ms")
attachedHash?.takeUnless(::isHashCurrent)?.let(::scheduleIdleDrop)
throw e
} catch (e: Exception) {
Log.w(TAG, "startStream[$generation]: failed after ${System.currentTimeMillis() - startAt}ms", e)
Log.w(TAG, "Failed to start P2P stream", e)
attachedHash?.takeUnless(::isHashCurrent)?.let(::scheduleIdleDrop)
if (isCurrentGeneration(generation)) {
_state.value = P2pStreamingState.Error(e.message ?: "Unknown torrent error")
@ -230,15 +226,15 @@ actual object P2pStreamingEngine {
}
actual fun stopStream() {
Log.d(TAG, "stopStream: detach active stream")
detachActiveStream()?.let(::scheduleIdleDrop)
}
actual fun shutdown() {
Log.d(TAG, "shutdown: stopping engine")
val hash = detachActiveStream()
val idleHashes = cancelScheduledIdleDrops()
val warmup = synchronized(lifecycleLock) {
warmupCooldownJob?.cancel()
warmupCooldownJob = null
val job = warmupJob
warmupJob = null
job
@ -275,12 +271,6 @@ actual object P2pStreamingEngine {
val preloadJob: Job?,
)
private data class ActiveFileSelection(
val requestedIdx: Int?,
val requestedFilename: String?,
val resolvedIdx: Int?,
)
private fun beginStreamGeneration(): DetachedStream {
val detached = synchronized(lifecycleLock) {
streamGeneration += 1
@ -311,6 +301,13 @@ actual object P2pStreamingEngine {
job?.join()
}
private fun cancelWarmupCooldown() {
synchronized(lifecycleLock) {
warmupCooldownJob?.cancel()
warmupCooldownJob = null
}
}
private fun attachTorrentIfCurrent(generation: Long, hash: String): Boolean =
synchronized(lifecycleLock) {
if (streamGeneration != generation) return@synchronized false
@ -333,7 +330,6 @@ actual object P2pStreamingEngine {
private fun scheduleIdleDrop(hash: String, delayMs: Long = IDLE_TORRENT_TTL_MS) {
val key = hashKey(hash)
if (key.isBlank()) return
Log.d(TAG, "idleDrop: scheduling hash=$hash in ${delayMs}ms")
val job = scope.launch {
delay(delayMs)
val shouldDrop = synchronized(lifecycleLock) {
@ -346,15 +342,11 @@ actual object P2pStreamingEngine {
}
}
if (shouldDrop) {
Log.d(TAG, "idleDrop: dropping hash=$hash after idle ttl")
api.dropTorrent(hash)
} else {
Log.d(TAG, "idleDrop: keeping hash=$hash because it is current again")
}
}
synchronized(lifecycleLock) {
if (hashMatches(currentHash, hash)) {
Log.d(TAG, "idleDrop: skip schedule for current hash=$hash")
job.cancel()
return
}
@ -366,7 +358,6 @@ actual object P2pStreamingEngine {
private fun cancelIdleDrop(hash: String) {
synchronized(lifecycleLock) {
idleDropJobs.remove(hashKey(hash))?.let {
Log.d(TAG, "idleDrop: cancelled hash=$hash")
it.cancel()
}
}
@ -456,124 +447,46 @@ actual object P2pStreamingEngine {
.trim()
.takeIf { it.isNotEmpty() }
private fun summarizeMagnet(magnetLink: String): String {
val query = magnetLink.substringAfter('?', missingDelimiterValue = "")
val params = query.split('&').filter { it.isNotBlank() }
val xt = params.firstOrNull { it.startsWith("xt=", ignoreCase = true) }
?.substringAfter('=')
.orEmpty()
val dn = params.firstOrNull { it.startsWith("dn=", ignoreCase = true) }
?.substringAfter('=')
?.let(::decodeQueryValue)
.orEmpty()
val trackerCount = params.count { it.startsWith("tr=", ignoreCase = true) }
return "xt=$xt dn=$dn trackers=$trackerCount length=${magnetLink.length}"
}
private fun summarizeStreamUrl(streamUrl: String): String {
val index = streamUrl.substringAfter("index=", missingDelimiterValue = "")
.substringBefore('&')
.takeIf { it.isNotBlank() }
?: "?"
val fileIdx = streamUrl.substringAfter("fileIdx=", missingDelimiterValue = "")
.substringBefore('&')
.takeIf { it.isNotBlank() }
?: "?"
val filename = streamUrl.substringAfter("filename=", missingDelimiterValue = "")
.substringBefore('&')
.takeIf { it.isNotBlank() }
?.let(::decodeQueryValue)
.orEmpty()
return "index=$index fileIdx=$fileIdx filename=$filename length=${streamUrl.length}"
}
private suspend fun resolveFileIndex(hash: String, requestedIdx: Int?, filename: String?): Int {
val requestedName = filename?.trim()?.takeIf { it.isNotEmpty() }
if (requestedIdx != null) {
val torrServerIndex = requestedIdx + 1
Log.d(
TAG,
"resolveFileIndex: requested fileIdx=$requestedIdx mapsToTorrServerId=$torrServerIndex " +
"filename=${requestedName.orEmpty()}",
)
if (requestedName != null) {
val files = waitForTorrentFiles(
hash = hash,
timeoutMs = FILE_INDEX_FAST_VALIDATION_TIMEOUT_MS,
reason = "fast-validation",
)
if (files.isNotEmpty()) {
logTorrentFiles(
label = "fast-validation",
files = files,
requestedIdx = requestedIdx,
)
resolveByFilename(files, requestedName)?.let { match ->
Log.d(
TAG,
"resolveFileIndex: filename match overrides requested index " +
"requestedIdx=$requestedIdx requestedTorrServerId=$torrServerIndex " +
"matchedId=${match.id} path=${match.path}",
)
return match.id
}
val requestedIdMatch = files.firstOrNull { it.id == torrServerIndex }
if (requestedIdMatch != null) {
Log.d(
TAG,
"resolveFileIndex: using requested TorrServer id=$torrServerIndex " +
"path=${requestedIdMatch.path} size=${formatBytes(requestedIdMatch.length)} " +
"because filename=$requestedName did not match metadata",
)
return torrServerIndex
}
val positionalFile = files.getOrNull(requestedIdx)
if (positionalFile != null) {
Log.w(
TAG,
"resolveFileIndex: TorrServer id=$torrServerIndex missing; using positional file " +
"position=$requestedIdx id=${positionalFile.id} path=${positionalFile.path}",
)
return positionalFile.id
}
Log.w(
TAG,
"resolveFileIndex: requestedIdx=$requestedIdx is outside metadata files=${files.size}; " +
"falling back to TorrServer id=$torrServerIndex",
)
} else {
Log.w(
TAG,
"resolveFileIndex: no metadata within ${FILE_INDEX_FAST_VALIDATION_TIMEOUT_MS}ms; " +
"using requested TorrServer id=$torrServerIndex",
)
}
} else {
Log.d(TAG, "resolveFileIndex: no filename to validate; using requested TorrServer id=$torrServerIndex")
}
return torrServerIndex
}
Log.d(TAG, "resolveFileIndex: no requested fileIdx; waiting for metadata filename=${requestedName.orEmpty()}")
val files = waitForTorrentFiles(
hash = hash,
timeoutMs = FILE_INDEX_METADATA_TIMEOUT_MS,
reason = "full-resolution",
)
if (files.isEmpty()) {
Log.w(TAG, "resolveFileIndex: no files after metadata timeout; guessing TorrServer id=1")
return 1
}
logTorrentFiles(label = "full-resolution", files = files, requestedIdx = requestedIdx)
if (requestedName != null) {
resolveByFilename(files, requestedName)?.let { match ->
Log.d(TAG, "resolveFileIndex: filename match id=${match.id} path=${match.path}")
return match.id
}
}
@ -586,19 +499,12 @@ actual object P2pStreamingEngine {
.maxByOrNull { it.length }
val result = videoFile?.id ?: files.maxByOrNull { it.length }?.id ?: 1
val resultFile = files.firstOrNull { it.id == result }
Log.d(
TAG,
"resolveFileIndex: largest video fallback id=$result " +
"path=${resultFile?.path.orEmpty()} size=${resultFile?.let { formatBytes(it.length) }.orEmpty()}",
)
return result
}
private suspend fun waitForTorrentFiles(
hash: String,
timeoutMs: Long,
reason: String,
): List<TorrServerFile> {
val startedAt = System.currentTimeMillis()
var attempt = 0
@ -607,18 +513,10 @@ actual object P2pStreamingEngine {
val stats = api.getTorrentStats(hash)
val files = stats?.files.orEmpty()
if (files.isNotEmpty()) {
Log.d(
TAG,
"metadata[$reason]: files=${files.size} after ${System.currentTimeMillis() - startedAt}ms " +
"attempt=$attempt peers=${stats?.peers ?: 0} seeds=${stats?.seeds ?: 0} " +
"download=${stats?.downloadSpeed ?: 0}Bps preload=${stats?.preloadedBytes ?: 0}",
)
return files
}
Log.d(TAG, "metadata[$reason]: waiting attempt=$attempt elapsed=${System.currentTimeMillis() - startedAt}ms")
delay(FILE_INDEX_POLL_INTERVAL_MS)
}
Log.w(TAG, "metadata[$reason]: unavailable after ${System.currentTimeMillis() - startedAt}ms")
return emptyList()
}
@ -628,7 +526,6 @@ actual object P2pStreamingEngine {
file.path.substringAfterLast('/').equals(name, ignoreCase = true)
}
if (exactBasename != null) {
Log.d(TAG, "resolveByFilename: exact basename match filename=$name id=${exactBasename.id}")
return exactBasename
}
@ -636,7 +533,6 @@ actual object P2pStreamingEngine {
file.path.equals(name, ignoreCase = true)
}
if (exactPath != null) {
Log.d(TAG, "resolveByFilename: exact path match filename=$name id=${exactPath.id}")
return exactPath
}
@ -644,98 +540,12 @@ actual object P2pStreamingEngine {
file.path.contains(name, ignoreCase = true)
}
if (contains != null) {
Log.d(TAG, "resolveByFilename: contains match filename=$name id=${contains.id}")
return contains
}
Log.w(TAG, "resolveByFilename: no match for filename=$name")
return null
}
private fun logTorrentFiles(
label: String,
files: List<TorrServerFile>,
requestedIdx: Int?,
) {
Log.d(TAG, "files[$label]: total=${files.size} requestedIdx=${requestedIdx ?: -1}")
val selectedPositions = buildSet {
addAll(0 until minOf(files.size, 25))
requestedIdx?.let { idx ->
add(idx)
add(idx - 2)
add(idx - 1)
add(idx + 1)
add(idx + 2)
}
files.indices.maxByOrNull { files[it].length }?.let(::add)
add(files.lastIndex)
}.filter { it in files.indices }.sorted()
selectedPositions.forEach { position ->
val file = files[position]
val requestedMarker = if (
position == requestedIdx ||
file.id == requestedIdx?.plus(1) ||
file.index == requestedIdx
) {
" requested"
} else {
""
}
Log.d(
TAG,
"files[$label]: pos=$position id=${file.id} index=${file.index}$requestedMarker " +
"size=${formatBytes(file.length)} path=${file.path}",
)
}
val omitted = files.size - selectedPositions.size
if (omitted > 0) {
Log.d(TAG, "files[$label]: omitted=$omitted middle entries")
}
}
private fun logStreamingFileAudit(
files: List<TorrServerFile>,
fileSelection: ActiveFileSelection,
) {
logTorrentFiles(
label = "streaming-audit",
files = files,
requestedIdx = fileSelection.requestedIdx,
)
val selected = fileSelection.resolvedIdx?.let { resolvedIdx ->
files.firstOrNull { it.id == resolvedIdx }
}
val filename = fileSelection.requestedFilename?.trim()?.takeIf { it.isNotEmpty() }
val filenameMatch = filename?.let { resolveByFilename(files, it) }
val rawIndexMatch = fileSelection.requestedIdx?.let { requestedIdx ->
files.firstOrNull { it.index == requestedIdx }
}
Log.d(
TAG,
"streaming-audit: resolvedId=${fileSelection.resolvedIdx} " +
"selectedPath=${selected?.path.orEmpty()} selectedSize=${selected?.let { formatBytes(it.length) }.orEmpty()} " +
"requestedIdx=${fileSelection.requestedIdx ?: -1} requestedFilename=${filename.orEmpty()} " +
"filenameMatchId=${filenameMatch?.id ?: -1} filenameMatchPath=${filenameMatch?.path.orEmpty()} " +
"rawIndexMatchId=${rawIndexMatch?.id ?: -1} rawIndexMatchPath=${rawIndexMatch?.path.orEmpty()}",
)
if (
fileSelection.resolvedIdx != null &&
filenameMatch != null &&
filenameMatch.id != fileSelection.resolvedIdx
) {
Log.w(
TAG,
"streaming-audit: FILE_SELECTION_MISMATCH resolvedId=${fileSelection.resolvedIdx} " +
"selectedPath=${selected?.path.orEmpty()} filenameMatchId=${filenameMatch.id} " +
"filenameMatchPath=${filenameMatch.path}",
)
}
}
private fun startPreload(
hash: String,
generation: Long,
@ -743,26 +553,15 @@ actual object P2pStreamingEngine {
selector: TorrServerStreamSelector,
) {
val job = scope.launch {
val preloadAt = System.currentTimeMillis()
Log.d(
TAG,
"preload[$generation]: starting hash=$hash legacyIndex=${selector.legacyIndex ?: -1} " +
"fileIdx=${selector.fileIdx ?: -1} filename=${selector.filename.orEmpty()}",
)
try {
val completed = api.preloadTorrent(
api.preloadTorrent(
magnetLink = magnetLink,
selector = selector,
)
Log.d(
TAG,
"preload[$generation]: completed=$completed in ${System.currentTimeMillis() - preloadAt}ms",
)
} catch (e: CancellationException) {
Log.d(TAG, "preload[$generation]: cancelled after ${System.currentTimeMillis() - preloadAt}ms")
throw e
} catch (e: Exception) {
Log.w(TAG, "preload[$generation]: failed after ${System.currentTimeMillis() - preloadAt}ms", e)
Log.w(TAG, "Torrent preload failed", e)
}
}
val previousJob = synchronized(lifecycleLock) {
@ -778,38 +577,15 @@ actual object P2pStreamingEngine {
previousJob?.cancel()
}
private fun formatBytes(bytes: Long): String =
when {
bytes >= 1_073_741_824L -> "${"%.2f".format(Locale.US, bytes / 1_073_741_824.0)}GB"
bytes >= 1_048_576L -> "${"%.1f".format(Locale.US, bytes / 1_048_576.0)}MB"
bytes >= 1_024L -> "${bytes / 1_024}KB"
else -> "${bytes}B"
}
private fun formatSpeed(bytesPerSecond: Long): String =
"${formatBytes(bytesPerSecond)}/s"
private fun startStatsPolling(
hash: String,
generation: Long,
fileSelection: ActiveFileSelection,
) {
val job = scope.launch {
var loggedMetadataAudit = false
while (isActive) {
if (!isCurrentGeneration(generation)) return@launch
try {
val stats = api.getTorrentStats(hash)
if (!loggedMetadataAudit && stats?.files?.isNotEmpty() == true) {
loggedMetadataAudit = true
logStreamingFileAudit(
files = stats.files,
fileSelection = fileSelection,
)
}
if (stats != null) {
Log.d(TAG, "stats[$generation]: ${stats.summary()}")
}
val currentState = _state.value
if (
stats != null &&
@ -845,16 +621,6 @@ actual object P2pStreamingEngine {
previousJob?.cancel()
}
private fun TorrServerStats.summary(): String =
"state=${status.ifBlank { "?" }} " +
"down=${formatSpeed(downloadSpeed)} up=${formatSpeed(uploadSpeed)} " +
"peers=$peers/$totalPeers seeds=$seeds pending=$pendingPeers halfOpen=$halfOpenPeers " +
"preload=${formatBytes(preloadedBytes)}/${formatBytes(preloadSize)} " +
"loaded=${formatBytes(loadedSize)}/${formatBytes(torrentSize)} " +
"read=${formatBytes(bytesRead)} useful=${formatBytes(bytesReadUsefulData)} " +
"chunks=$chunksRead/$chunksReadUseful wasted=$chunksReadWasted " +
"piecesGood=$piecesDirtiedGood piecesBad=$piecesDirtiedBad"
private val DEFAULT_TRACKERS = listOf(
"udp://tracker.opentrackr.org:1337/announce",
"udp://open.demonii.com:1337/announce",
@ -896,7 +662,6 @@ actual object P2pStreamingEngine {
suspend fun start() = startMutex.withLock {
withContext(Dispatchers.IO) {
if (isRunning()) {
Log.d(TAG, "TorrServer already running")
return@withContext
}
@ -923,15 +688,12 @@ actual object P2pStreamingEngine {
processBuilder.directory(configDir)
processBuilder.redirectErrorStream(true)
Log.d(TAG, "Starting TorrServer on port $PORT from ${binaryFile.absolutePath}")
process = processBuilder.start()
val proc = process!!
Thread {
try {
proc.inputStream.bufferedReader().forEachLine { line ->
Log.d(TAG, "[server] $line")
}
proc.inputStream.bufferedReader().forEachLine { }
} catch (_: Exception) {
}
}.apply {
@ -942,7 +704,6 @@ actual object P2pStreamingEngine {
val deadline = System.currentTimeMillis() + STARTUP_TIMEOUT_MS
while (System.currentTimeMillis() < deadline) {
if (isRunning()) {
Log.d(TAG, "TorrServer started successfully")
return@withContext
}
if (!isProcessAlive(process)) {
@ -985,7 +746,6 @@ actual object P2pStreamingEngine {
}
}
process = null
Log.d(TAG, "TorrServer stopped")
}
private fun killOrphanedProcess() {
@ -993,7 +753,6 @@ actual object P2pStreamingEngine {
val request = Request.Builder().url("$baseUrl/shutdown").build()
healthClient.newCall(request).execute().close()
Thread.sleep(1_000L)
Log.d(TAG, "Shut down orphaned TorrServer instance")
} catch (_: Exception) {
}
}
@ -1082,7 +841,6 @@ actual object P2pStreamingEngine {
success = false,
changed = false,
)
val beforeSummary = summarizeSettings(settings)
val changes = mutableListOf<String>()
putIfDifferent(settings, "CacheSize", STREAMING_CACHE_SIZE_BYTES, changes)
@ -1112,17 +870,12 @@ actual object P2pStreamingEngine {
putIfDifferent(settings, "StoreSettingsInJson", true, changes)
if (changes.isEmpty()) {
Log.d(TAG, "streaming-settings: already tuned $beforeSummary")
return@withContext StreamingSettingsResult(
success = true,
changed = false,
)
}
Log.d(
TAG,
"streaming-settings: applying ${changes.joinToString()} from $beforeSummary",
)
val body = JSONObject().apply {
put("action", "set")
put("sets", settings)
@ -1141,7 +894,6 @@ actual object P2pStreamingEngine {
)
}
}
Log.d(TAG, "streaming-settings: applied ${summarizeSettings(settings)}")
StreamingSettingsResult(
success = true,
changed = true,
@ -1185,7 +937,7 @@ actual object P2pStreamingEngine {
val currentValue = settings.optInt(key, Int.MIN_VALUE)
if (!hasKey || currentValue != desiredValue) {
settings.put(key, desiredValue)
changes += "$key=${formatSettingsValue(currentValue, hasKey)}->$desiredValue"
changes += key
}
}
@ -1199,7 +951,7 @@ actual object P2pStreamingEngine {
val currentValue = settings.optLong(key, Long.MIN_VALUE)
if (!hasKey || currentValue != desiredValue) {
settings.put(key, desiredValue)
changes += "$key=${formatSettingsValue(currentValue, hasKey)}->$desiredValue"
changes += key
}
}
@ -1213,35 +965,10 @@ actual object P2pStreamingEngine {
val currentValue = settings.optBoolean(key, !desiredValue)
if (!hasKey || currentValue != desiredValue) {
settings.put(key, desiredValue)
changes += "$key=${if (hasKey) currentValue else "unset"}->$desiredValue"
changes += key
}
}
private fun formatSettingsValue(value: Int, hasKey: Boolean): String =
if (hasKey) value.toString() else "unset"
private fun formatSettingsValue(value: Long, hasKey: Boolean): String =
if (hasKey) value.toString() else "unset"
private fun summarizeSettings(settings: JSONObject): String =
"cache=${formatBytes(settings.optLong("CacheSize", 0))} " +
"connections=${settings.optInt("ConnectionsLimit", 0)} " +
"halfOpen=${settings.optInt("HalfOpenConnectionsLimit", 0)}/" +
"${settings.optInt("TotalHalfOpenConnectionsLimit", 0)} " +
"peerWater=${settings.optInt("TorrentPeersLowWater", 0)}/" +
"${settings.optInt("TorrentPeersHighWater", 0)} " +
"dialMs=${settings.optInt("MinDialTimeoutMs", 0)}/" +
"${settings.optInt("NominalDialTimeoutMs", 0)} " +
"handshakeMs=${settings.optInt("HandshakeTimeoutMs", 0)} " +
"readAhead=${settings.optInt("ReaderReadAHead", 0)} " +
"preload=${settings.optInt("PreloadCache", 0)} " +
"responsive=${settings.optBoolean("ResponsiveMode", false)} " +
"dht=${!settings.optBoolean("DisableDHT", false)} " +
"pex=${!settings.optBoolean("DisablePEX", false)} " +
"tcp=${!settings.optBoolean("DisableTCP", false)} " +
"utp=${!settings.optBoolean("DisableUTP", false)} " +
"rate=${settings.optInt("DownloadRateLimit", 0)}/${settings.optInt("UploadRateLimit", 0)}KBps"
suspend fun addTorrent(magnetLink: String, title: String? = null): String? = withContext(Dispatchers.IO) {
val body = JSONObject().apply {
put("action", "add")
@ -1263,7 +990,6 @@ actual object P2pStreamingEngine {
}
val json = JSONObject(response.body?.string() ?: "{}")
val hash = json.optString("hash", "")
Log.d(TAG, "Torrent added: $hash")
hash.ifEmpty { null }
}
} catch (e: Exception) {
@ -1281,7 +1007,6 @@ actual object P2pStreamingEngine {
selector = selector,
mode = "preload",
)
Log.d(TAG, "preload: url=${summarizeStreamUrl(url)}")
val request = Request.Builder()
.url(url)
.get()
@ -1377,7 +1102,6 @@ actual object P2pStreamingEngine {
try {
client.newCall(request).execute().close()
Log.d(TAG, "Torrent dropped: $hash")
} catch (e: Exception) {
Log.w(TAG, "dropTorrent error", e)
}

View file

@ -45,9 +45,7 @@ object P2pSettingsRepository {
if (p2pEnabled == enabled) return
p2pEnabled = enabled
P2pSettingsStorage.saveP2pEnabled(enabled)
if (enabled) {
P2pStreamingEngine.warmup()
} else {
if (!enabled) {
P2pStreamingEngine.shutdown()
}
publish()
@ -126,6 +124,7 @@ class P2pStreamingException(message: String) : Exception(message)
expect object P2pStreamingEngine {
val state: StateFlow<P2pStreamingState>
fun warmup()
fun cooldownWarmup()
suspend fun startStream(request: P2pStreamRequest): String
fun stopStream()
fun shutdown()

View file

@ -53,6 +53,7 @@ import androidx.compose.material3.Icon
import androidx.compose.material3.MaterialTheme
import androidx.compose.material3.Text
import androidx.compose.runtime.Composable
import androidx.compose.runtime.DisposableEffect
import androidx.compose.runtime.LaunchedEffect
import androidx.compose.runtime.getValue
import androidx.compose.runtime.mutableStateOf
@ -91,6 +92,8 @@ import coil3.compose.AsyncImage
import com.nuvio.app.core.ui.nuvioSafeBottomPadding
import com.nuvio.app.features.debrid.DebridProviders
import com.nuvio.app.features.debrid.DebridSettingsRepository
import com.nuvio.app.features.p2p.P2pSettingsRepository
import com.nuvio.app.features.p2p.P2pStreamingEngine
import com.nuvio.app.features.player.PlayerSettingsRepository
import com.nuvio.app.features.watchprogress.WatchProgressRepository
import kotlinx.coroutines.launch
@ -140,6 +143,10 @@ fun StreamsScreen(
DebridSettingsRepository.ensureLoaded()
DebridSettingsRepository.uiState
}.collectAsStateWithLifecycle()
val p2pSettings by remember {
P2pSettingsRepository.ensureLoaded()
P2pSettingsRepository.uiState
}.collectAsStateWithLifecycle()
val watchProgressUiState by remember {
WatchProgressRepository.ensureLoaded()
WatchProgressRepository.uiState
@ -181,6 +188,15 @@ fun StreamsScreen(
}
}
DisposableEffect(P2pSettingsRepository.isVisible, p2pSettings.p2pEnabled) {
if (P2pSettingsRepository.isVisible && p2pSettings.p2pEnabled) {
P2pStreamingEngine.warmup()
}
onDispose {
P2pStreamingEngine.cooldownWarmup()
}
}
LaunchedEffect(type, videoId, seasonNumber, episodeNumber, manualSelection) {
StreamsRepository.load(
type = type,

View file

@ -12,6 +12,10 @@ actual object P2pStreamingEngine {
_state.value = P2pStreamingState.Idle
}
actual fun cooldownWarmup() {
_state.value = P2pStreamingState.Idle
}
actual suspend fun startStream(request: P2pStreamRequest): String {
val message = "P2P streaming is not available on this platform"
_state.value = P2pStreamingState.Error(message)