diff --git a/composeApp/src/androidMain/jniLibs/arm64-v8a/libtorrserver.so b/composeApp/src/androidMain/jniLibs/arm64-v8a/libtorrserver.so index a2c2b4d17..844e339ed 100755 Binary files a/composeApp/src/androidMain/jniLibs/arm64-v8a/libtorrserver.so and b/composeApp/src/androidMain/jniLibs/arm64-v8a/libtorrserver.so differ diff --git a/composeApp/src/androidMain/jniLibs/armeabi-v7a/libtorrserver.so b/composeApp/src/androidMain/jniLibs/armeabi-v7a/libtorrserver.so index 1e6f2580e..3c9935edd 100755 Binary files a/composeApp/src/androidMain/jniLibs/armeabi-v7a/libtorrserver.so and b/composeApp/src/androidMain/jniLibs/armeabi-v7a/libtorrserver.so differ diff --git a/composeApp/src/androidMain/jniLibs/x86/libtorrserver.so b/composeApp/src/androidMain/jniLibs/x86/libtorrserver.so index 9c3c2a6c9..edc2bf081 100755 Binary files a/composeApp/src/androidMain/jniLibs/x86/libtorrserver.so and b/composeApp/src/androidMain/jniLibs/x86/libtorrserver.so differ diff --git a/composeApp/src/androidMain/jniLibs/x86_64/libtorrserver.so b/composeApp/src/androidMain/jniLibs/x86_64/libtorrserver.so index bebe25f34..5fc0b7a13 100755 Binary files a/composeApp/src/androidMain/jniLibs/x86_64/libtorrserver.so and b/composeApp/src/androidMain/jniLibs/x86_64/libtorrserver.so differ diff --git a/composeApp/src/androidMain/kotlin/com/nuvio/app/features/p2p/P2pStreamingEngine.android.kt b/composeApp/src/androidMain/kotlin/com/nuvio/app/features/p2p/P2pStreamingEngine.android.kt index 641ad8c0c..d9e952b9e 100644 --- a/composeApp/src/androidMain/kotlin/com/nuvio/app/features/p2p/P2pStreamingEngine.android.kt +++ b/composeApp/src/androidMain/kotlin/com/nuvio/app/features/p2p/P2pStreamingEngine.android.kt @@ -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 { 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, - 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, - 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() 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) } diff --git a/composeApp/src/commonMain/kotlin/com/nuvio/app/features/p2p/P2pStreaming.kt b/composeApp/src/commonMain/kotlin/com/nuvio/app/features/p2p/P2pStreaming.kt index a037adea1..9a2c5c516 100644 --- a/composeApp/src/commonMain/kotlin/com/nuvio/app/features/p2p/P2pStreaming.kt +++ b/composeApp/src/commonMain/kotlin/com/nuvio/app/features/p2p/P2pStreaming.kt @@ -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 fun warmup() + fun cooldownWarmup() suspend fun startStream(request: P2pStreamRequest): String fun stopStream() fun shutdown() diff --git a/composeApp/src/commonMain/kotlin/com/nuvio/app/features/streams/StreamsScreen.kt b/composeApp/src/commonMain/kotlin/com/nuvio/app/features/streams/StreamsScreen.kt index b107d575c..0c3385e24 100644 --- a/composeApp/src/commonMain/kotlin/com/nuvio/app/features/streams/StreamsScreen.kt +++ b/composeApp/src/commonMain/kotlin/com/nuvio/app/features/streams/StreamsScreen.kt @@ -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, diff --git a/composeApp/src/iosMain/kotlin/com/nuvio/app/features/p2p/P2pStreamingEngine.ios.kt b/composeApp/src/iosMain/kotlin/com/nuvio/app/features/p2p/P2pStreamingEngine.ios.kt index 43a2fde95..ecca96a16 100644 --- a/composeApp/src/iosMain/kotlin/com/nuvio/app/features/p2p/P2pStreamingEngine.ios.kt +++ b/composeApp/src/iosMain/kotlin/com/nuvio/app/features/p2p/P2pStreamingEngine.ios.kt @@ -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)