From 74bd4bb09adcdb8e4635d7fd5aba74d764e63cdc Mon Sep 17 00:00:00 2001 From: Lin Zhang Date: Sun, 20 Sep 2026 15:50:34 +0800 Subject: [PATCH] Torrent v2 8/8: V2 engine lifecycle, rate controls and restart recovery MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Consolidates #227–#228, preserving source b55c48167d84b452e1abfb40f8e6f778f2f348f8. Includes the foundation listener cancellation regression fix discovered by consolidation CI. --- docs/design/torrent-v2-integrity.md | 82 ++++ .../ketch/torrent/KotlinTorrentEngine.kt | 139 ++++++- .../ketch/torrent/TorrentRateLimiter.kt | 45 +- .../ketch/torrent/TorrentSessionAdmission.kt | 37 ++ .../ketch/torrent/TorrentV2DownloadSession.kt | 178 ++++++++ .../ketch/torrent/TorrentV2PieceStore.kt | 13 + .../ketch/torrent/TorrentV2SessionLoop.kt | 33 +- .../ketch/torrent/KotlinTorrentV2OwnerTest.kt | 385 ++++++++++++++++++ .../ketch/torrent/TorrentRateLimiterTest.kt | 30 +- .../torrent/TorrentV2DownloadSessionTest.kt | 239 +++++++++++ .../ketch/torrent/TorrentV2EngineRateTest.kt | 90 ++++ .../ketch/torrent/TorrentV2SessionLoopTest.kt | 42 ++ 12 files changed, 1285 insertions(+), 28 deletions(-) create mode 100644 library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentV2DownloadSession.kt create mode 100644 library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/KotlinTorrentV2OwnerTest.kt create mode 100644 library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentV2DownloadSessionTest.kt create mode 100644 library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentV2EngineRateTest.kt diff --git a/docs/design/torrent-v2-integrity.md b/docs/design/torrent-v2-integrity.md index 67eb09f14..25273ccff 100644 --- a/docs/design/torrent-v2-integrity.md +++ b/docs/design/torrent-v2-integrity.md @@ -1064,3 +1064,85 @@ The metadata exchange partition has a minimum of `TRACKER_SCRAPE_WORKSPACE_BYTES of the accepted metainfo-size limit. Small valid metainfo limits therefore retain scrape capability. Configuration validation includes this floor in the aggregate exchange ceiling; an explicitly undersized aggregate budget fails at construction instead of silently disabling every scrape. + +## V2 download lifecycle ownership + +`TorrentV2DownloadSession` composes the verified full-metainfo pipeline under one scoped owner. +Resume initializes/rechecks owned storage before discovery and resets published verification +progress during checking. Each run owns the endpoint producer, bounded dialer, peer pool, and +commit worker; terminal completion is published only after those scopes finish cleanup. + +Pause joins the active run, including discovery and provider writes, before acknowledging the +paused state. Owner exit also joins the lifetime and closes storage before returning admission. +Ordinary failures publish a stopped state and can be retried; immediate retries join the prior +terminal job. Metadata identity and normalized selection must match the store before any I/O. +Progress notifications use committed/rechecked storage bytes, not bytes received from peers. + +Real TCP coverage downloads a v2 file, alters its payload while paused, and verifies that resume +rechecks and repairs it before reporting completion. Additional tests cover delayed discovery +cleanup, repeated failure/retry, store-selection rejection before filesystem creation, and owner +shutdown admission ordering. Both TCP endpoints are Ketch fixtures, not independent v2 interop. + +The engine must still admit document/layout/store indexes before constructing this owner and keep +that admission until it returns. Engine registration, full-identity incoming routing, public source +v2 resolution, checkpoints/TaskStore wiring, rate controls, seeding, and tracker/public discovery +integration remain required follow-up work. This owner does not advertise public v2 support. + +Lifecycle admission also compares the complete layout to storage's canonical hybrid-aware layout: +file IDs/indices, offsets and lengths, piece length, and protocol/payload totals. Matching only the +full info hash is insufficient because omitting hybrid mapping changes IDs after padding entries. +A regression now rejects that mismatch before I/O while accepting the canonical selected-file ID. + +## Engine-owned v2 download lifetimes + +`KotlinTorrentEngine.withV2Download` runs the full-metainfo lifecycle as an engine-owned child. +The engine admits retained content and storage/layout indexes before construction, enforces the +configured file/piece/info limits, and uses its shared connection, transfer, session, and payload +handle pools. Caller cancellation joins the owned child; engine shutdown cancels and joins it +before the registered metadata admission is released. No output files are deleted by owner exit. + +V1 and v2 registrations share the active-task ceiling and canonical output-overlap checks. Full +v2 hashes identify registrations; a hybrid's v1 identity cannot also own a legacy session, in +either registration order. Legacy removal cannot release a live v2 output claim. Failed admission +or registration returns its memory and leaves no output ownership behind. Hybrid layouts are +constructed with their authenticated v1 padding/index mapping, including selected files after gaps. + +The scoped engine path takes policy-authorized endpoints from its caller. Tracker/DHT integration, +incoming full-identity routing, source/SDK v2 resolution, checkpoint/TaskStore restore, rate controls, +and seeding remain pending. It is not advertised by the public engine/source API. Real TCP tests +exercise both a pure v2 file and a hybrid selected file after padding; both peers are Ketch fixtures. + +Engine shutdown uses one shared cleanup operation outside the engine job tree. Calls from an +engine-owned callback request that operation without joining their own ancestor; external callers +await the same operation as the full cleanup barrier. Concurrent close/stop calls cannot create +multiple cleanup owners, and start rejects a runtime whose shutdown has been requested. Callback +regressions cover both the scoped v2 body and its discovery producer, followed by an external +stop that proves all registration admission has returned. + +## V2 runtime capability integration (in progress) + +Future changes are grouped around a usable v2 runtime capability instead of one PR per helper. +The local integration branch now applies global and task download budgets before submitting block +requests. Admission checks both buckets atomically; a blocked bucket or rejected command queue +consumes neither. Rate waits do not block peer readers or the session's event handling. Retry +delays follow available credit and are capped at 50 ms to observe live limit changes. + +Rate accounting covers requested payload, including requests that are sent but later fail; protocol +overhead is separate. The existing 16 KiB burst remains. Real TCP and scheduler tests verify that +global and task limits both apply and can be removed while downloading, plus transactional queue +rejection and retry delays that do not impose a fixed polling throughput ceiling. Broader v2 +discovery, recovery, source/API integration and production gates still need work before this +capability branch is ready for a PR. + +The scoped runtime also accepts a previously decoded v2 checkpoint. Engine admission includes its +retained ownership records, path strings and hint bitmap before store construction. Session setup +validates and adopts ownership before exposing the owner; failure closes the store and releases +registration and memory. The first resume rehashes actual payloads before reporting progress or +completion. Later resumes recheck again without reapplying the initial snapshot. A checkpoint +cannot authorize replacement files or a different task, selection, content identity or destination. + +Recovery tests persist and reload the checkpoint and authenticated catalog, then enter the engine +with fresh storage ownership. They cover completion without discovery, payload modification while +paused, invalid task binding and admission cleanup. Decoding/catalog I/O are still caller-owned; +this entry point does not yet bound their transient allocations or schedule automatic checkpoints. +TaskStore/source wiring and durable checkpoint publication remain part of the pending integration. diff --git a/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/KotlinTorrentEngine.kt b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/KotlinTorrentEngine.kt index 9bda888df..6fa37c4ed 100644 --- a/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/KotlinTorrentEngine.kt +++ b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/KotlinTorrentEngine.kt @@ -4,12 +4,14 @@ import kotlinx.coroutines.CancellationException import kotlinx.coroutines.CoroutineStart import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.Deferred import kotlinx.coroutines.Job import kotlinx.coroutines.NonCancellable import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.async import kotlinx.coroutines.awaitAll import kotlinx.coroutines.cancel +import kotlinx.coroutines.cancelAndJoin import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.channels.SendChannel import kotlinx.coroutines.coroutineScope @@ -26,7 +28,12 @@ import kotlinx.coroutines.withContext import kotlinx.coroutines.withTimeout import kotlinx.coroutines.withTimeoutOrNull import okio.FileSystem +import okio.Path +import okio.ByteString.Companion.toByteString import okio.Path.Companion.toPath +import kotlin.coroutines.AbstractCoroutineContextElement +import kotlin.coroutines.CoroutineContext +import kotlin.concurrent.atomics.AtomicReference import kotlin.concurrent.atomics.AtomicBoolean import kotlin.concurrent.atomics.ExperimentalAtomicApi @@ -41,7 +48,13 @@ internal class KotlinTorrentEngine( private val discoveryIntervalMs: Long = 30_000, private val nowMs: () -> Long = monotonicClock(), ) : TorrentEngine { - private val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default) + private class RuntimeContext : AbstractCoroutineContextElement(Key) { + companion object Key : CoroutineContext.Key + } + + private val runtimeContext = RuntimeContext() + private val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default + runtimeContext) + private val shutdown = AtomicReference?>(null) private val network = TorrentConnectionBudget(rawNetwork, config.maxConnections) private val exchangeBudgets = TorrentExchangeBudgets(config) private val budget = exchangeBudgets.transfer @@ -55,6 +68,7 @@ internal class KotlinTorrentEngine( private val mutex = Mutex() private val dhtMutex = Mutex() private val sessions = mutableMapOf() + private val v2Identities = mutableMapOf() private val outputs = mutableMapOf() private val sessionLeases = mutableMapOf() private val running = AtomicBoolean(false) @@ -71,7 +85,7 @@ internal class KotlinTorrentEngine( } override suspend fun start() = mutex.withLock { - check(!closed) { "Torrent runtime is closed" } + check(!closed && shutdown.load() == null) { "Torrent runtime is closed" } if (running.load()) return@withLock val listener = network.listen(PeerEndpoint("0.0.0.0", config.listenPort)) port = listener.local.port @@ -121,14 +135,31 @@ internal class KotlinTorrentEngine( } finally { listener.close() } } - override fun close() { - running.store(false) - scope.cancel() - network.close() - http.close() - } + override fun close() { requestShutdown() } + /** Internal callbacks request shutdown; external callers await the complete cleanup barrier. */ override suspend fun stop() { + val pending = requestShutdown() + if (currentCoroutineContext()[RuntimeContext] === runtimeContext) return + withContext(NonCancellable) { pending.await() } + } + + private fun requestShutdown(): Deferred { + shutdown.load()?.let { return it } + // This job cannot be a descendant of the engine it must cancel and join. + val candidate = CoroutineScope(Dispatchers.Default).async(start = CoroutineStart.LAZY) { + stopRuntime() + } + if (shutdown.compareAndSet(null, candidate)) { + running.store(false) + candidate.start() + return candidate + } + candidate.cancel() + return checkNotNull(shutdown.load()) + } + + private suspend fun stopRuntime() { val active = mutex.withLock { if (closed) return closed = true @@ -298,9 +329,13 @@ internal class KotlinTorrentEngine( override suspend fun addTask(spec: TorrentTaskSpec): KotlinTorrentSession = mutex.withLock { check(isRunning && !closed) - check(sessions.size < config.maxActiveTorrents) { "Too many active torrents" } + check(sessions.size + v2Identities.size < config.maxActiveTorrents) { + "Too many active torrents" + } val hash = spec.metadata.infoHash.hex - check(hash !in sessions) { "Torrent already has an active owner" } + check(hash !in sessions && v2Identities.values.none { it.v1 == spec.metadata.infoHash }) { + "Torrent already has an active owner" + } val lease = admissions.admit(spec, config) try { val requested = FileSystem.SYSTEM.canonicalize(".".toPath()) @@ -308,13 +343,7 @@ internal class KotlinTorrentEngine( val store = TorrentPieceStore(spec.metadata, requested, spec.selected, spec.taskId, storageSlots = storageSlots) val output = store.outputPath.toPath() - check(outputs.values.none { previous -> - val path = previous.toPath() - val left = path.segments.map { canonicalTorrentName(it).lowercase() } - val right = output.segments.map { canonicalTorrentName(it).lowercase() } - path.root.toString().equals(output.root.toString(), ignoreCase = true) && - (left.take(right.size) == right || right.take(left.size) == left) - }) { "Torrent output overlaps another task" } + requireAvailableOutput(output) val checkpoint = spec.resumeData?.let(TorrentCheckpoint::decode) val trackerState = TorrentBufferBudget( maxOf(1, trackerControlStateWeight().toInt())) @@ -340,6 +369,79 @@ internal class KotlinTorrentEngine( } } + private fun requireAvailableOutput(output: Path) { + check(outputs.values.none { previous -> + val path = previous.toPath() + val left = path.segments.map { canonicalTorrentName(it).lowercase() } + val right = output.segments.map { canonicalTorrentName(it).lowercase() } + path.root.toString().equals(output.root.toString(), ignoreCase = true) && + (left.take(right.size) == right || right.take(left.size) == left) + }) { "Torrent output overlaps another task" } + } + + /** + * Scoped full-metainfo download integration. The caller supplies policy-authorized endpoints; + * this path does not yet register incoming v2 routes, resolve magnets, or seed after completion. + * Retained metadata/storage indexes are admitted before storage construction or file I/O. + * The caller owns checkpoint decoding admission; retained recovery records are admitted here. + * Recovery validates ownership before exposing the session and rechecks payloads on resume. + */ + suspend fun withV2Download( + taskId: String, + document: TorrentV2Document, + outputPath: String, + selected: Set = emptySet(), + checkpoint: TorrentV2Checkpoint? = null, + discover: suspend (SendChannel) -> Unit, + body: suspend (TorrentV2DownloadSession) -> T, + ): T { + val pending = scope.async { + val hash = document.info.hash.hex + var lease: TorrentBufferBudget.Lease? = null + var registered = false + try { + val prepared = mutex.withLock { + check(isRunning && !closed) { "Torrent runtime is closed" } + check(sessions.size + v2Identities.size < config.maxActiveTorrents) { + "Too many active torrents" + } + check(hash !in v2Identities && + (document.identity.v1?.hex !in sessions) && + v2Identities.values.none { identity -> + identity.v1 != null && identity.v1 == document.identity.v1 + }) { "Torrent already has an active owner" } + lease = admitV2Session(document, selected.size, outputPath.length, config, + exchangeBudgets.sessions, checkpoint) + val selection = selected.toSet() + val requested = FileSystem.SYSTEM.canonicalize(".".toPath()) + .resolve(outputPath).normalized() + val parent = checkNotNull(requested.parent) { "Output root must have a parent" } + val output = FileSystem.SYSTEM.canonicalize(parent) / requested.name + requireAvailableOutput(output) + val layout = TorrentContentLayout.from(document.info, document.hybrid) + val store = TorrentV2PieceStore(document, output, selection, taskId, budget, storageSlots) + v2Identities[hash] = document.identity + outputs[hash] = output.toString() + registered = true + Triple(layout, store, selection) + } + TorrentV2DownloadSession.run(document, prepared.first, prepared.third, prepared.second, + network, peerId.toByteString(), budget, exchangeBudgets.sessions, + maxPeers = minOf(config.connectionsPerTorrent, 500), + globalRate = downloadRate, checkpoint = checkpoint, discover = discover, body = body) + } finally { + withContext(NonCancellable) { + try { + if (registered) mutex.withLock { v2Identities.remove(hash); outputs.remove(hash) } + } finally { lease?.close() } + } + } + } + try { return pending.await() } finally { + withContext(NonCancellable) { pending.cancelAndJoin() } + } + } + private suspend fun discover( spec: TorrentTaskSpec, output: SendChannel, @@ -455,7 +557,8 @@ internal class KotlinTorrentEngine( override suspend fun removeTorrent(infoHash: String, deleteFiles: Boolean) { mutex.withLock { - sessions[infoHash]?.close(deleteFiles) + val session = sessions[infoHash] ?: return@withLock + session.close(deleteFiles) sessions.remove(infoHash) outputs.remove(infoHash) sessionLeases.remove(infoHash)?.let(admissions::release) diff --git a/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentRateLimiter.kt b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentRateLimiter.kt index c1c596502..0d08f65a7 100644 --- a/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentRateLimiter.kt +++ b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentRateLimiter.kt @@ -5,6 +5,7 @@ import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.withLock import kotlin.concurrent.atomics.AtomicLong import kotlin.concurrent.atomics.ExperimentalAtomicApi +import kotlin.math.ceil /** Live rate changes take effect within 50 ms, including removing a limit while waiting. */ @OptIn(ExperimentalAtomicApi::class) @@ -24,6 +25,46 @@ internal class TorrentRateLimiter( rate.store(bytesPerSecond) } + /** + * Nonblocking block-request admission. Zero consumes both buckets; a positive delay consumes + * neither. Queue rejection also consumes neither. The admission callback must not suspend or + * reenter a limiter. Always invoke on the engine/global bucket with a distinct task-local bucket so lock + * ordering remains global then task. Waiting sessions poll live rate changes within 50 ms. + */ + suspend fun requestDelay( + bytes: Int, + task: TorrentRateLimiter, + admit: () -> Boolean = { true }, + ): Long { + require(bytes in 1..16_384 && task !== this) + return mutex.withLock { + task.mutex.withLock { + val globalRate = rate.load() + val taskRate = task.rate.load() + refill(globalRate) + task.refill(taskRate) + val delay = maxOf(delayFor(bytes, globalRate), task.delayFor(bytes, taskRate)) + if (delay == 0L && admit()) { + if (globalRate != 0L) tokens -= bytes + if (taskRate != 0L) task.tokens -= bytes + } + delay + } + } + } + + private fun refill(currentRate: Long) { + val now = nowMs() + tokens = minOf(16_384.0, + tokens + (now - updated).coerceAtLeast(0) * currentRate.toDouble() / 1000.0) + updated = now + } + + private fun delayFor(bytes: Int, currentRate: Long): Long { + if (currentRate == 0L || tokens >= bytes) return 0 + return ceil((bytes - tokens) * 1000 / currentRate).toLong().coerceIn(1, 50) + } + suspend fun acquire(bytes: Int) { require(bytes >= 0) var remaining = bytes @@ -31,9 +72,7 @@ internal class TorrentRateLimiter( val consumed = mutex.withLock { val currentRate = rate.load() if (currentRate == 0L) return - val now = nowMs() - tokens = minOf(16_384.0, tokens + (now - updated).coerceAtLeast(0) * currentRate.toDouble() / 1000.0) - updated = now + refill(currentRate) minOf(remaining, tokens.toInt()).also { tokens -= it } } remaining -= consumed diff --git a/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentSessionAdmission.kt b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentSessionAdmission.kt index be65b1697..ef5309a9d 100644 --- a/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentSessionAdmission.kt +++ b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentSessionAdmission.kt @@ -79,3 +79,40 @@ internal class TorrentAdmissionLedger(private val budget: TorrentBufferBudget) { entries.exchange(null)?.forEach { it.close() } } } + +/** Admit retained full v2 content plus storage/layout indexes before constructing their arrays. */ +internal fun admitV2Session( + document: TorrentV2Document, + selectedCount: Int, + outputLength: Int, + config: TorrentConfig, + budget: TorrentBufferBudget, + checkpoint: TorrentV2Checkpoint? = null, +): TorrentBufferBudget.Lease { + require(document.info.rawInfo.size <= config.maxMetadataBytes) { "Torrent info limit exceeded" } + require(outputLength >= 0) + require(document.info.files.size <= config.maxFilesPerTorrent && + selectedCount in 0..config.maxFilesPerTorrent) { "Torrent file limit exceeded" } + require(document.info.pieceLength <= 16 * 1024 * 1024) { "Torrent piece length limit exceeded" } + var pieces = 0L + for (file in document.info.files) { + val count = file.length / document.info.pieceLength + + if (file.length % document.info.pieceLength == 0L) 0 else 1 + require(count <= config.maxPiecesPerTorrent - pieces) { "Torrent piece limit exceeded" } + pieces += count + } + val recoveryBytes = checkpoint?.let { saved -> + 4096L + saved.output.length * 4L + saved.verifiedHint.size * 2L + + saved.selected.sumOf { it.length * 4L + 128 } + + saved.owned.sumOf { claim -> + 512L + claim.identity.length * 4L + + claim.components.sumOf { it.length * 4L + 128 } + } + } ?: 0L + val bytes = recoveryBytes + document.info.rawInfo.size * 4L + + document.pieceLayers.values.sumOf { it.size * 2L + 128 } + + document.info.files.sumOf { file -> 2048L + file.path.sumOf { it.size * 8L + 128 } } + + pieces * 128 + selectedCount * 128L + outputLength * 4L + 256 * 1024 + check(bytes <= budget.capacity) { "Torrent session state exceeds admission capacity" } + return checkNotNull(budget.reserve(bytes.toInt())) { "Torrent session state budget exhausted" } +} diff --git a/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentV2DownloadSession.kt b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentV2DownloadSession.kt new file mode 100644 index 000000000..0ea495663 --- /dev/null +++ b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentV2DownloadSession.kt @@ -0,0 +1,178 @@ +package com.linroid.ketch.torrent + +import kotlinx.coroutines.CancellationException +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Job +import kotlinx.coroutines.NonCancellable +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.cancelAndJoin +import kotlinx.coroutines.channels.Channel +import kotlinx.coroutines.channels.SendChannel +import kotlinx.coroutines.coroutineScope +import kotlinx.coroutines.currentCoroutineContext +import kotlinx.coroutines.ensureActive +import kotlinx.coroutines.flow.MutableStateFlow +import kotlinx.coroutines.flow.StateFlow +import kotlinx.coroutines.launch +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock +import kotlinx.coroutines.withContext +import okio.ByteString + +/** Full-metainfo v2 download owner. Network policy and metadata admission belong to the engine. */ +internal class TorrentV2DownloadSession private constructor( + private val scope: CoroutineScope, + private val document: TorrentV2Document, + private val layout: TorrentContentLayout, + private val selected: Set, + private val store: TorrentV2PieceStore, + private val network: TorrentNetwork, + private val peerId: ByteString, + private val buffers: TorrentBufferBudget, + private val memory: TorrentBufferBudget, + private val maxPeers: Int, + private val globalRate: TorrentRateLimiter, + private val discover: suspend (SendChannel) -> Unit, +) { + private val lifecycle = Mutex() + private val downloadRate = TorrentRateLimiter() + private var job: Job? = null + private var closed = false + private val mutableState = MutableStateFlow(TorrentSessionState.PAUSED) + private val mutableProgress = MutableStateFlow(0L) + private val mutableFailure = MutableStateFlow(null) + val state: StateFlow get() = mutableState + val verifiedBytes: StateFlow get() = mutableProgress + val failure: StateFlow get() = mutableFailure + + fun setDownloadRateLimit(bytesPerSecond: Long) = downloadRate.set(bytesPerSecond) + + suspend fun resume() = lifecycle.withLock { + check(!closed) { "Torrent session is closed" } + currentCoroutineContext().ensureActive() + if (job?.isActive == true && (mutableState.value == TorrentSessionState.CHECKING_FILES || + mutableState.value == TorrentSessionState.DOWNLOADING)) return@withLock + job?.join() + checkNotNull(scope.coroutineContext[Job]).ensureActive() + mutableProgress.value = 0 + mutableFailure.value = null + mutableState.value = TorrentSessionState.CHECKING_FILES + job = scope.launch { + try { + store.initialize() + // Persisted bits and earlier live progress cannot authorize bytes changed while paused. + store.recheck() + updateProgress() + if (!store.completed()) { + mutableState.value = TorrentSessionState.DOWNLOADING + transfer() + } + currentCoroutineContext().ensureActive() + mutableState.value = TorrentSessionState.FINISHED + } catch (error: CancellationException) { + if (!checkNotNull(currentCoroutineContext()[Job]).isActive) throw error + mutableFailure.value = error + mutableState.value = TorrentSessionState.STOPPED + } catch (error: Exception) { + mutableFailure.value = error + mutableState.value = TorrentSessionState.STOPPED + } + } + } + + suspend fun pause() = lifecycle.withLock { + check(!closed) { "Torrent session is closed" } + withContext(NonCancellable) { + stop() + mutableState.value = TorrentSessionState.PAUSED + } + } + + private suspend fun stop() = withContext(NonCancellable) { + job?.cancelAndJoin() + job = null + updateProgress() + } + + private suspend fun updateProgress() { mutableProgress.value = store.progress().values.sum() } + + private suspend fun transfer() = coroutineScope { + val endpoints = Channel(maxPeers) + val discovery = launch { + try { discover(endpoints) } catch (error: Throwable) { + endpoints.close(error) + } finally { endpoints.close() } + } + try { + PeerV2Pool.run(memory, maxPeers = maxPeers) { pool -> + PeerV2Dialer.run(endpoints, memory, parallelism = minOf(4, maxPeers), + connect = { remote -> + PeerV2Connector.connect(network, remote, document, layout, peerId, buffers, memory) + }, + ) { dialer -> + TorrentV2CommitWorker.run(store) { worker -> + TorrentV2SessionLoop.download(layout, selected, store, pool, worker, buffers, memory, + maxPeers = maxPeers, connections = dialer.connections, onProgress = ::updateProgress, + requestDelay = { bytes, admit -> globalRate.requestDelay(bytes, downloadRate, admit) }) + } + } + } + } finally { + withContext(NonCancellable) { + try { discovery.cancelAndJoin() } finally { endpoints.cancel() } + } + } + } + + private suspend fun shutdown() = lifecycle.withLock { + closed = true + try { stop() } finally { + try { checkNotNull(scope.coroutineContext[Job]).cancelAndJoin() } finally { + try { store.close() } finally { mutableState.value = TorrentSessionState.STOPPED } + } + } + } + + companion object { + /** + * Owns the store until all commands, discovery, peers and provider writes have joined. + * The caller admits document/layout/storage indexes and any checkpoint before entry, retaining + * that admission until return. This lifecycle does not yet seed or expose the public engine API. + */ + suspend fun run( + document: TorrentV2Document, + layout: TorrentContentLayout, + selected: Set, + store: TorrentV2PieceStore, + network: TorrentNetwork, + peerId: ByteString, + buffers: TorrentBufferBudget, + state: TorrentBufferBudget, + maxPeers: Int = 100, + globalRate: TorrentRateLimiter = TorrentRateLimiter(), + checkpoint: TorrentV2Checkpoint? = null, + discover: suspend (SendChannel) -> Unit, + body: suspend (TorrentV2DownloadSession) -> T, + ): T = coroutineScope { + require(layout.infoHash == document.info.hash && peerId.size == 20) + require(maxPeers in 1..500) + store.requireBinding(document.identity, selected, layout) + val lease = checkNotNull(state.reserve(maxPeers * 512 + 4096)) { + "Session lifecycle state budget exhausted" + } + val owner = SupervisorJob(coroutineContext[Job]) + val session = TorrentV2DownloadSession(CoroutineScope(coroutineContext + owner), document, + layout, selected.toSet(), store, network, peerId, buffers, state, maxPeers, + globalRate, discover) + try { + // Restore before exposing the owner. Resume always rechecks the adopted payloads. + if (checkpoint != null) store.restore(checkpoint) + body(session) + } finally { + withContext(NonCancellable) { + try { session.shutdown() } finally { lease.close() } + } + } + } + } +} diff --git a/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentV2PieceStore.kt b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentV2PieceStore.kt index a53ecd988..09224da29 100644 --- a/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentV2PieceStore.kt +++ b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentV2PieceStore.kt @@ -60,6 +60,19 @@ internal class TorrentV2PieceStore( verified = BooleanArray(verifier.layout.pieceCount.toInt()) } + /** Validate session wiring before starting any filesystem or network operation. */ + fun requireBinding(identity: TorrentIdentity, ids: Set, layout: TorrentContentLayout) { + require(document.identity == identity) { "Store belongs to another torrent" } + val canonical = verifier.layout + require(layout.infoHash == canonical.infoHash && layout.pieceLength == canonical.pieceLength && + layout.protocolBytes == canonical.protocolBytes && + layout.payloadBytes == canonical.payloadBytes && + layout.files == canonical.files) { "Session layout differs from storage layout" } + require(if (ids.isEmpty()) selected.size == mapping.files.size else ids == selected) { + "Store selection differs from session selection" + } + } + suspend fun initialize() = mutex.withLock { check(!closed) if (initialized) return@withLock diff --git a/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentV2SessionLoop.kt b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentV2SessionLoop.kt index ec96d6f50..ac47c468c 100644 --- a/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentV2SessionLoop.kt +++ b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentV2SessionLoop.kt @@ -15,6 +15,7 @@ internal object TorrentV2SessionLoop { var interestPending = false var inFlight = 0 var admissionBlocked = false + var rateRetryAt = 0L } private sealed interface Event { @@ -43,6 +44,9 @@ internal object TorrentV2SessionLoop { maxActive: Int = 2, pipeline: Int = 32, connections: ReceiveChannel? = null, + onProgress: suspend () -> Unit = {}, + requestDelay: suspend (Int, () -> Boolean) -> Long = { _, admit -> admit(); 0 }, + nowMs: () -> Long = monotonicClock(), ) { require(maxPeers in 1..500 && pipeline in 1..256) val lease = checkNotNull(state.reserve(maxPeers * 512 + 1024)) { @@ -61,10 +65,11 @@ internal object TorrentV2SessionLoop { fun wakeAdmission() { peers.values.forEach { it.admissionBlocked = false } } - fun pump() { + suspend fun pump() { pieces.submitReady(worker) for ((peer, view) in peers) { - if (view.admissionBlocked) continue + if (view.rateRetryAt <= nowMs()) view.rateRetryAt = 0 + if (view.admissionBlocked || view.rateRetryAt != 0L) continue val interested = view.needed > 0 if (interested != view.interested && !view.interestPending) { if (view.commands.trySend(PeerV2DownloadActor.Command.Interest(interested)).isSuccess) { @@ -82,10 +87,20 @@ internal object TorrentV2SessionLoop { } } if (plan != null) { - if (view.commands.trySend(PeerV2DownloadActor.Command.Request(plan)).isSuccess) { + var sent = false + val request = plan + val delay = requestDelay(request.request.length) { + view.commands.trySend(PeerV2DownloadActor.Command.Request(request)).isSuccess + .also { sent = it } + } + require(delay in 0..50 && (delay == 0L || !sent)) { "Invalid request rate delay" } + if (delay > 0) { + check(pieces.resolve(request, null)) + view.rateRetryAt = nowMs() + delay + } else if (sent) { view.inFlight++ } else { - check(pieces.resolve(plan, null)) + check(pieces.resolve(request, null)) view.admissionBlocked = true } } @@ -140,6 +155,7 @@ internal object TorrentV2SessionLoop { } } + onProgress() var preference = 0 var acceptingConnections = connections != null while (!pieces.completed()) { @@ -157,8 +173,12 @@ internal object TorrentV2SessionLoop { } } } - // Only admission pressure enables a retry timer; healthy idle peers do not poll. - if (peers.values.any { it.admissionBlocked }) onTimeout(250) { Event.Retry } + // Retry only blocked work; healthy idle peers do not poll. + val memoryDelay = if (peers.values.any { it.admissionBlocked }) 250L else null + val rateDelay = peers.values.filter { it.rateRetryAt != 0L } + .minOfOrNull { (it.rateRetryAt - nowMs()).coerceAtLeast(1) } + val retryDelay = listOfNotNull(memoryDelay, rateDelay).minOrNull() + if (retryDelay != null) onTimeout(retryDelay) { Event.Retry } } preference = (preference + 1) % 3 when (event) { @@ -181,6 +201,7 @@ internal object TorrentV2SessionLoop { view.needed-- } } + onProgress() wakeAdmission() } is Event.Retry -> wakeAdmission() diff --git a/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/KotlinTorrentV2OwnerTest.kt b/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/KotlinTorrentV2OwnerTest.kt new file mode 100644 index 000000000..e973ba56a --- /dev/null +++ b/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/KotlinTorrentV2OwnerTest.kt @@ -0,0 +1,385 @@ +package com.linroid.ketch.torrent + +import kotlinx.coroutines.CancellationException +import kotlinx.coroutines.currentCoroutineContext +import kotlinx.coroutines.ensureActive +import kotlinx.coroutines.CompletableDeferred +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.NonCancellable +import kotlinx.coroutines.async +import kotlinx.coroutines.awaitCancellation +import kotlinx.coroutines.cancelAndJoin +import kotlinx.coroutines.flow.first +import kotlinx.coroutines.launch +import kotlinx.coroutines.sync.Semaphore +import kotlinx.coroutines.test.runTest +import kotlinx.coroutines.withContext +import kotlinx.coroutines.withTimeout +import okio.ByteString.Companion.toByteString +import okio.FileSystem +import kotlin.test.Test +import kotlin.test.assertContentEquals +import kotlin.test.assertEquals +import kotlin.test.assertFailsWith +import kotlin.test.assertFalse +import kotlin.test.assertNotNull +import kotlin.test.assertIs +import kotlin.test.assertTrue + +class KotlinTorrentV2OwnerTest { + private val bytes = byteArrayOf(1, 2, 3, 4) + private fun document(name: String = "pack", hybrid: Boolean = false): TorrentV2Document { + val info = mutableMapOf("name" to name, "meta version" to 2L, + "piece length" to 16_384L, "file tree" to mapOf("a" to mapOf("" to mapOf( + "length" to 4L, "pieces root" to sha256Digest(bytes))))) + if (hybrid) { + info["files"] = listOf(mapOf("path" to listOf("a"), "length" to 4L)) + info["pieces"] = sha1Digest(bytes) + } + return TorrentV2Document.parse(Bencode.encode(mapOf("info" to info, + "piece layers" to emptyMap()))) + } + + private fun output() = FileSystem.SYSTEM_TEMPORARY_DIRECTORY / + "ketch-engine-v2-${InfoHash.fromBytes(torrentRandomBytes(20)).hex}" + + @Test + fun restartRechecksCheckpointBeforeCompletionAndAgainAfterPause() = runTest { + withContext(Dispatchers.Default) { + withTimeout(15_000) { + val root = output() + torrentFileSystem.createDirectory(root) + val destination = root / "payload" + val document = document() + val buffers = TorrentBufferBudget(1024 * 1024) + val original = TorrentV2PieceStore(document, destination, emptySet(), "restart", + buffers, Semaphore(1)) + val engine = KotlinTorrentEngine(TorrentConfig(dhtEnabled = false)) + try { + original.initialize() + assertTrue(original.commit(0, bytes)) + val catalog = TorrentContentCatalog(root / "catalog") + val saved = original.checkpoint(catalog) + TorrentCheckpointFile(root / "state", "restart", catalog).save(document, saved) + original.close() + val loaded = assertNotNull(TorrentCheckpointFile(root / "state", "restart", + TorrentContentCatalog(root / "catalog")).load()) + val checkpoint = loaded.checkpoint + val admission = TorrentBufferBudget(4 * 1024 * 1024) + val config = TorrentConfig(dhtEnabled = false) + val fresh = admitV2Session(document, 0, destination.toString().length, config, admission) + val freshBytes = admission.allocated + fresh.close() + val recovery = admitV2Session(document, 0, destination.toString().length, config, + admission, checkpoint) + assertTrue(admission.allocated > freshBytes) + recovery.close() + val tight = TorrentBufferBudget(freshBytes) + assertFailsWith { + admitV2Session(document, 0, destination.toString().length, config, tight, checkpoint) + } + assertEquals(0, tight.allocated) + engine.start() + val discovery = CompletableDeferred() + engine.withV2Download("restart", loaded.document, destination.toString(), + checkpoint = checkpoint, discover = { + discovery.complete(Unit) + awaitCancellation() + }) { session -> + assertEquals(0L, session.verifiedBytes.value) + session.resume() + session.state.first { it == TorrentSessionState.FINISHED } + assertEquals(4L, session.verifiedBytes.value) + assertFalse(discovery.isCompleted) + session.pause() + torrentFileSystem.write(destination / "a") { write(ByteArray(4)) } + session.resume() + discovery.await() + assertEquals(0L, session.verifiedBytes.value) + session.pause() + } + assertEquals(0, engine.admittedSessionBytes) + } finally { + engine.stop() + original.close() + torrentFileSystem.deleteRecursively(root, mustExist = false) + } + assertEquals(0, buffers.allocated) + } + } + } + + @Test + fun rejectedCheckpointReleasesEngineClaimsWithoutExposingTheSession() = runTest { + val root = output() + torrentFileSystem.createDirectory(root) + val destination = root / "payload" + val document = document() + val buffers = TorrentBufferBudget(1024 * 1024) + val original = TorrentV2PieceStore(document, destination, emptySet(), "restart", + buffers, Semaphore(1)) + val engine = KotlinTorrentEngine(TorrentConfig(dhtEnabled = false)) + try { + original.initialize() + assertTrue(original.commit(0, bytes)) + val checkpoint = original.checkpoint(TorrentContentCatalog(root / "catalog")) + original.close() + engine.start() + assertFailsWith { + engine.withV2Download("wrong-task", document, destination.toString(), + checkpoint = checkpoint, discover = { error("Unexpected discovery") }) { + error("Invalid recovery exposed the session") + } + } + assertEquals(0, engine.admittedSessionBytes) + engine.withV2Download("restart", document, destination.toString(), + checkpoint = checkpoint, discover = { error("Unexpected discovery") }) { session -> + assertEquals(0L, session.verifiedBytes.value) + } + assertContentEquals(bytes, torrentFileSystem.read(destination / "a") { readByteArray() }) + assertEquals(0, engine.admittedSessionBytes) + } finally { + engine.stop() + original.close() + torrentFileSystem.deleteRecursively(root, mustExist = false) + } + assertEquals(0, buffers.allocated) + } + + @Test + fun engineDownloadsV2ThroughSharedNetworkAndReleasesOwnershipAfterCompletion() = runTest { + transfer(hybrid = false) + } + + @Test + fun engineDownloadsSelectedHybridFileAfterPadding() = runTest { transfer(hybrid = true) } + + private suspend fun transfer(hybrid: Boolean) { + withContext(Dispatchers.Default) { + withTimeout(15_000) { + val document = if (!hybrid) document() else TorrentV2Document.parse(Bencode.encode(mapOf( + "info" to mapOf("name" to "hybrid", "meta version" to 2L, "piece length" to 16_384L, + "file tree" to mapOf( + "a" to mapOf("" to mapOf("length" to 4L, "pieces root" to sha256Digest(bytes))), + "b" to mapOf("" to mapOf("length" to 4L, "pieces root" to sha256Digest(bytes)))), + "files" to listOf(mapOf("path" to listOf("a"), "length" to 4L), + mapOf("attr" to "p", "length" to 16_380L), + mapOf("path" to listOf("b"), "length" to 4L)), + "pieces" to (sha1Digest(bytes + ByteArray(16_380)) + sha1Digest(bytes))), + "piece layers" to emptyMap() + ))) + val output = output() + val remote = createTorrentNetwork() + val listener = remote.listen(PeerEndpoint("127.0.0.1", 0)) + val buffers = TorrentBufferBudget(1024 * 1024) + val engine = KotlinTorrentEngine(TorrentConfig(dhtEnabled = false)) + val server = async { + val connection = listener.accept() + try { + PeerIdentityHandshake(document.identity).respond(connection, + ByteArray(20) { 2 }.toByteString(), buffers) + val wire = PeerWire(connection, pieceCount = if (hybrid) 2 else 1) + wire.send(PeerMessage.Bitfield(byteArrayOf(if (hybrid) 64 else 128.toByte()))) + wire.send(PeerMessage.Control(PeerMessage.Signal.UNCHOKE)) + assertEquals(PeerMessage.Control(PeerMessage.Signal.INTERESTED), wire.read()) + val request = assertIs(wire.read()) + assertEquals(if (hybrid) 1 else 0, request.index) + wire.send(PeerMessage.Piece(request.index, request.begin, bytes)) + } finally { connection.close() } + } + try { + engine.start() + engine.withV2Download("v2", document, output.toString(), + selected = if (hybrid) setOf("2") else emptySet(), + discover = { it.send(listener.local); awaitCancellation() }) { session -> + assertTrue(engine.admittedSessionBytes > 0) + session.resume() + session.state.first { it == TorrentSessionState.FINISHED } + assertEquals(4L, session.verifiedBytes.value) + val file = output / if (hybrid) "b" else "a" + assertContentEquals(bytes, torrentFileSystem.read(file) { readByteArray() }) + if (hybrid) assertFalse(torrentFileSystem.exists(output / "a")) + } + assertEquals(0, engine.admittedSessionBytes) + server.await() + } finally { + engine.stop() + listener.close() + remote.close() + server.cancelAndJoin() + torrentFileSystem.deleteRecursively(output, mustExist = false) + } + assertEquals(0, buffers.allocated) + } + } + } + + @Test + fun engineShutdownJoinsV2DiscoveryBeforeReturningAdmission() = runTest { + withContext(Dispatchers.Default) { + withTimeout(15_000) { + val engine = KotlinTorrentEngine(TorrentConfig(dhtEnabled = false)) + val output = output() + val entered = CompletableDeferred() + val cleaning = CompletableDeferred() + val release = CompletableDeferred() + engine.start() + val owner = launch { + engine.withV2Download("v2", document(), output.toString(), discover = { + entered.complete(Unit) + try { awaitCancellation() } finally { + withContext(NonCancellable) { cleaning.complete(Unit); release.await() } + } + }) { it.resume(); awaitCancellation() } + } + try { + entered.await() + val stop = async { engine.stop() } + cleaning.await() + assertFalse(stop.isCompleted) + assertTrue(engine.admittedSessionBytes > 0) + release.complete(Unit) + stop.await() + owner.join() + assertEquals(0, engine.admittedSessionBytes) + } finally { + release.complete(Unit) + owner.cancelAndJoin() + engine.stop() + torrentFileSystem.deleteRecursively(output, mustExist = false) + } + } + } + } + + @Test + fun callbacksCanRequestShutdownWithoutJoiningTheirOwnJob() = runTest { + withContext(Dispatchers.Default) { + withTimeout(15_000) { + for (fromDiscovery in listOf(false, true)) { + val engine = KotlinTorrentEngine(TorrentConfig(dhtEnabled = false)) + val output = output() + val returned = CompletableDeferred() + engine.start() + try { + try { + engine.withV2Download("v2", document(), output.toString(), discover = { + engine.stop() + returned.complete(Unit) + }) { session -> + if (fromDiscovery) { + session.resume() + awaitCancellation() + } else { + engine.stop() + returned.complete(Unit) + } + } + } catch (_: CancellationException) { + currentCoroutineContext().ensureActive() + } + returned.await() + engine.stop() + assertFalse(engine.isRunning) + assertEquals(0, engine.admittedSessionBytes) + assertFailsWith { engine.start() } + } finally { + engine.stop() + torrentFileSystem.deleteRecursively(output, mustExist = false) + } + } + } + } + } + + @Test + fun duplicateIdentityAndOverlappingOutputsRejectBeforeDiscovery() = runTest { + val engine = KotlinTorrentEngine(TorrentConfig(dhtEnabled = false)) + val output = output() + try { + engine.start() + engine.withV2Download("first", document(), output.toString(), discover = {}) { + val retained = engine.admittedSessionBytes + // A legacy removal request cannot release the scoped v2 owner's output claim. + engine.removeTorrent(document().info.hash.hex, deleteFiles = false) + assertFailsWith { + engine.withV2Download("same", document(), (output.parent!! / "unused").toString(), + discover = { error("Unexpected discovery") }) { error("Unexpected owner") } + } + assertFailsWith { + engine.withV2Download("overlap", document("other"), output.toString(), + discover = { error("Unexpected discovery") }) { error("Unexpected owner") } + } + assertEquals(retained, engine.admittedSessionBytes) + assertFalse(torrentFileSystem.exists(output)) + } + engine.withV2Download("retry", document(), output.toString(), discover = {}) {} + assertEquals(0, engine.admittedSessionBytes) + } finally { engine.stop() } + } + + @Test + fun hybridIdentityCannotHaveIndependentV1AndV2Owners() = runTest { + val engine = KotlinTorrentEngine(TorrentConfig(dhtEnabled = false)) + val hybrid = document(hybrid = true) + // The engine accepts trusted metadata adapters; model a v1 adapter for the same raw info. + val legacy = TorrentMetadata(checkNotNull(hybrid.identity.v1), "a", 16_384, 4, + listOf(TorrentMetadata.TorrentFile(0, "a", 4)), + infoBytes = hybrid.info.rawInfo.toByteArray(), pieceHashes = sha1Digest(bytes)) + val spec = TorrentTaskSpec("legacy", legacy, output().toString(), emptySet()) + try { + engine.start() + engine.withV2Download("v2", hybrid, output().toString(), discover = {}) { + assertFailsWith { engine.addTask(spec) } + } + engine.addTask(spec) + val retained = engine.admittedSessionBytes + assertFailsWith { + engine.withV2Download("v2", hybrid, output().toString(), discover = {}) { + error("Unexpected owner") + } + } + assertEquals(retained, engine.admittedSessionBytes) + engine.removeTorrent(legacy.infoHash.hex, deleteFiles = false) + assertEquals(0, engine.admittedSessionBytes) + } finally { engine.stop() } + } + + @Test + fun activeTaskLimitIsSharedAcrossV1AndV2() = runTest { + val engine = KotlinTorrentEngine(TorrentConfig(dhtEnabled = false, maxActiveTorrents = 1)) + val legacy = TorrentMetadata.fromBencode(Bencode.encode(mapOf("info" to mapOf( + "name" to "legacy", "length" to 4L, "piece length" to 4L, "pieces" to sha1Digest(bytes) + )))) + val spec = TorrentTaskSpec("legacy", legacy, output().toString(), emptySet()) + try { + engine.start() + engine.withV2Download("v2", document(), output().toString(), discover = {}) { + assertFailsWith { engine.addTask(spec) } + } + engine.addTask(spec) + assertFailsWith { + engine.withV2Download("v2", document(), output().toString(), discover = {}) { + error("Unexpected owner") + } + } + } finally { engine.stop() } + assertEquals(0, engine.admittedSessionBytes) + } + + @Test + fun metadataAdmissionFailsBeforeStorageAndDoesNotLeakOwnership() = runTest { + val engine = KotlinTorrentEngine(TorrentConfig(dhtEnabled = false, maxSessionStateBytes = 1)) + val output = output() + try { + engine.start() + assertFailsWith { + engine.withV2Download("v2", document(), output.toString(), discover = {}) { + error("Unexpected owner") + } + } + assertFalse(torrentFileSystem.exists(output)) + assertEquals(0, engine.admittedSessionBytes) + } finally { engine.stop() } + } +} diff --git a/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentRateLimiterTest.kt b/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentRateLimiterTest.kt index a38543028..365bf5b3c 100644 --- a/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentRateLimiterTest.kt +++ b/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentRateLimiterTest.kt @@ -5,11 +5,39 @@ import kotlinx.coroutines.test.advanceTimeBy import kotlinx.coroutines.test.runCurrent import kotlinx.coroutines.test.runTest import kotlin.test.Test +import kotlin.test.assertEquals import kotlin.test.assertFalse import kotlin.test.assertTrue +@OptIn(kotlinx.coroutines.ExperimentalCoroutinesApi::class) class TorrentRateLimiterTest { - @OptIn(kotlinx.coroutines.ExperimentalCoroutinesApi::class) + @Test + fun rejectedQueueAndBlockedBucketDoNotConsumeTheOtherBucket() = runTest { + val global = TorrentRateLimiter(1) { testScheduler.currentTime } + val task = TorrentRateLimiter(1) { testScheduler.currentTime } + assertEquals(0L, global.requestDelay(16_384, task) { false }) + assertEquals(0L, global.requestDelay(16_384, task)) + val freshTask = TorrentRateLimiter(1) { testScheduler.currentTime } + assertEquals(50L, global.requestDelay(16_384, freshTask)) + val freshGlobal = TorrentRateLimiter(1) { testScheduler.currentTime } + assertEquals(50L, freshGlobal.requestDelay(16_384, task)) + assertEquals(0L, freshGlobal.requestDelay(16_384, freshTask)) + } + + @Test + fun requestRetryDelayTracksRateWithoutAnArtificialPollingThroughputCeiling() = runTest { + val global = TorrentRateLimiter(1024 * 1024) { testScheduler.currentTime } + val task = TorrentRateLimiter() { testScheduler.currentTime } + assertEquals(0L, global.requestDelay(16_384, task)) + assertEquals(16L, global.requestDelay(16_384, task)) + advanceTimeBy(16) + assertEquals(0L, global.requestDelay(16_384, task)) + global.set(1) + assertEquals(50L, global.requestDelay(16_384, task)) + global.set(0) + assertEquals(0L, global.requestDelay(16_384, task)) + } + @Test fun limit_enforcesSustainedRateAndRemovingLimitUnblocksWaiters() = runTest { val limiter = TorrentRateLimiter(1024) { testScheduler.currentTime } diff --git a/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentV2DownloadSessionTest.kt b/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentV2DownloadSessionTest.kt new file mode 100644 index 000000000..645fee994 --- /dev/null +++ b/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentV2DownloadSessionTest.kt @@ -0,0 +1,239 @@ +package com.linroid.ketch.torrent + +import kotlinx.coroutines.CompletableDeferred +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.NonCancellable +import kotlinx.coroutines.async +import kotlinx.coroutines.awaitCancellation +import kotlinx.coroutines.cancelAndJoin +import kotlinx.coroutines.flow.first +import kotlinx.coroutines.launch +import kotlinx.coroutines.sync.Semaphore +import kotlinx.coroutines.test.runTest +import kotlinx.coroutines.withContext +import kotlinx.coroutines.withTimeout +import okio.ByteString.Companion.toByteString +import okio.FileSystem +import okio.IOException +import kotlin.test.Test +import kotlin.test.assertContentEquals +import kotlin.test.assertEquals +import kotlin.test.assertFalse +import kotlin.test.assertFailsWith +import kotlin.test.assertIs +import kotlin.test.assertTrue + +class TorrentV2DownloadSessionTest { + private val bytes = byteArrayOf(1, 2, 3, 4) + private val document = TorrentV2Document.parse(Bencode.encode(mapOf( + "info" to mapOf("name" to "pack", "meta version" to 2L, "piece length" to 16_384L, + "file tree" to mapOf("a" to mapOf("" to mapOf("length" to 4L, + "pieces root" to sha256Digest(bytes))))), + "piece layers" to emptyMap() + ))) + private val layout = TorrentContentLayout.from(document.info) + private val peerId = ByteArray(20) { 1 }.toByteString() + + private inner class Fixture(val scope: CoroutineScope) { + val output = FileSystem.SYSTEM_TEMPORARY_DIRECTORY / + "ketch-v2-owner-${InfoHash.fromBytes(torrentRandomBytes(20)).hex}" + val buffers = TorrentBufferBudget(4 * 1024 * 1024) + val state = TorrentBufferBudget(4 * 1024 * 1024) + val network = createTorrentNetwork() + val store = TorrentV2PieceStore(document, output, emptySet(), "owner", buffers, Semaphore(4)) + + suspend fun run( + discover: suspend (kotlinx.coroutines.channels.SendChannel) -> Unit, + body: suspend (TorrentV2DownloadSession) -> T, + ): T = TorrentV2DownloadSession.run(document, layout, emptySet(), store, network, peerId, + buffers, state, maxPeers = 1, discover = discover, body = body) + } + + private suspend fun fixture(body: suspend Fixture.() -> Unit) = withContext(Dispatchers.Default) { + withTimeout(15_000) { + val fixture = Fixture(this) + try { fixture.body() } finally { + withContext(NonCancellable) { + fixture.network.close() + fixture.store.cleanup() + } + } + assertEquals(0, fixture.buffers.allocated) + assertEquals(0, fixture.state.allocated) + } + } + + @Test + fun resumedDownloadRechecksModifiedPayloadAndPublishesOnlyVerifiedBytes() = runTest { + fixture { + val listener = network.listen(PeerEndpoint("127.0.0.1", 0)) + val secondRequest = CompletableDeferred() + val release = CompletableDeferred() + val server = scope.async { + repeat(2) { round -> + val connection = listener.accept() + try { + PeerIdentityHandshake(document.identity).respond(connection, + ByteArray(20) { 2 }.toByteString(), buffers) + val wire = PeerWire(connection, pieceCount = 1) + wire.send(PeerMessage.Bitfield(byteArrayOf(128.toByte()))) + wire.send(PeerMessage.Control(PeerMessage.Signal.UNCHOKE)) + assertEquals(PeerMessage.Control(PeerMessage.Signal.INTERESTED), wire.read()) + val request = assertIs(wire.read()) + if (round == 1) { secondRequest.complete(Unit); release.await() } + wire.send(PeerMessage.Piece(request.index, request.begin, bytes)) + } finally { connection.close() } + } + } + try { + run(discover = { it.send(listener.local); awaitCancellation() }) { session -> + session.resume() + session.state.first { it == TorrentSessionState.FINISHED } + assertEquals(4L, session.verifiedBytes.value) + session.pause() + torrentFileSystem.write(output / "a") { write(ByteArray(4)) } + session.resume() + secondRequest.await() + assertEquals(0L, session.verifiedBytes.value) + release.complete(Unit) + session.state.first { it == TorrentSessionState.FINISHED } + assertEquals(4L, session.verifiedBytes.value) + assertContentEquals(bytes, torrentFileSystem.read(output / "a") { readByteArray() }) + } + server.await() + } finally { listener.close(); server.cancelAndJoin() } + } + } + + @Test + fun pauseWaitsForDiscoveryCleanupAndResumeCreatesANewLifetime() = runTest { + fixture { + val entered = CompletableDeferred() + val cleaning = CompletableDeferred() + val release = CompletableDeferred() + val restarted = CompletableDeferred() + var starts = 0 + run(discover = { + starts++ + if (starts == 2) restarted.complete(Unit) + entered.complete(Unit) + try { awaitCancellation() } finally { + withContext(NonCancellable) { cleaning.complete(Unit); release.await() } + } + }) { session -> + val retained = state.allocated + session.resume() + entered.await() + val pause = scope.launch { session.pause() } + cleaning.await() + assertFalse(pause.isCompleted) + assertTrue(state.allocated >= retained) + release.complete(Unit) + pause.join() + assertEquals(TorrentSessionState.PAUSED, session.state.value) + assertEquals(retained, state.allocated) + assertEquals(0, buffers.allocated) + session.resume() + restarted.await() + session.pause() + assertEquals(2, starts) + } + } + } + + @Test + fun discoveryFailureStopsSessionAndCanBeRetried() = runTest { + fixture { + var attempts = 0 + run(discover = { attempts++; throw IOException("Discovery failed") }) { session -> + repeat(2) { + session.resume() + session.state.first { it == TorrentSessionState.STOPPED } + assertIs(session.failure.value) + assertEquals(0L, session.verifiedBytes.value) + } + assertEquals(2, attempts) + } + } + } + + @Test + fun ownerExitJoinsDiscoveryBeforeReturningAdmissionAndClosingTheStore() = runTest { + fixture { + val ready = CompletableDeferred() + val entered = CompletableDeferred() + val finish = CompletableDeferred() + val cleaning = CompletableDeferred() + val release = CompletableDeferred() + val owner = scope.async { + run(discover = { + entered.complete(Unit) + try { awaitCancellation() } finally { + withContext(NonCancellable) { cleaning.complete(Unit); release.await() } + } + }) { session -> ready.complete(session); finish.await() } + } + val session = ready.await() + session.resume() + entered.await() + finish.complete(Unit) + cleaning.await() + assertFalse(owner.isCompleted) + assertTrue(state.allocated > 0) + release.complete(Unit) + owner.await() + assertEquals(0, state.allocated) + assertEquals(0, buffers.allocated) + assertEquals(TorrentSessionState.STOPPED, session.state.value) + assertFailsWith { session.resume() } + assertFailsWith { store.initialize() } + } + } + + @Test + fun rejectsHybridLayoutThatChangesSelectedV1FileIdsBeforeIo() = runTest { + val hybrid = TorrentV2Document.parse(Bencode.encode(mapOf("info" to mapOf( + "name" to "hybrid", "meta version" to 2L, "piece length" to 16_384L, + "file tree" to mapOf( + "a" to mapOf("" to mapOf("length" to 4L, "pieces root" to sha256Digest(bytes))), + "b" to mapOf("" to mapOf("length" to 4L, "pieces root" to sha256Digest(bytes)))), + "files" to listOf(mapOf("path" to listOf("a"), "length" to 4L), + mapOf("attr" to "p", "length" to 16_380L), + mapOf("path" to listOf("b"), "length" to 4L)), + "pieces" to (sha1Digest(bytes + ByteArray(16_380)) + sha1Digest(bytes)) + ), "piece layers" to emptyMap()))) + val output = FileSystem.SYSTEM_TEMPORARY_DIRECTORY / + "ketch-layout-bind-${InfoHash.fromBytes(torrentRandomBytes(20)).hex}" + val buffers = TorrentBufferBudget(1024 * 1024) + val state = TorrentBufferBudget(1024 * 1024) + val network = createTorrentNetwork() + val store = TorrentV2PieceStore(hybrid, output, setOf("2"), "hybrid", buffers, Semaphore(1)) + try { + assertFailsWith { + TorrentV2DownloadSession.run(hybrid, TorrentContentLayout.from(hybrid.info), setOf("2"), + store, network, peerId, buffers, state, discover = { error("Unexpected discovery") }) { + error("Accepted a layout with the wrong file IDs") + } + } + assertFalse(torrentFileSystem.exists(output)) + assertEquals(0, state.allocated) + TorrentV2DownloadSession.run(hybrid, TorrentContentLayout.from(hybrid.info, hybrid.hybrid), + setOf("2"), store, network, peerId, buffers, state, discover = {}) {} + } finally { network.close(); store.cleanup() } + assertEquals(0, state.allocated) + } + + @Test + fun selectionMismatchRejectsBeforeCreatingFilesOrStartingDiscovery() = runTest { + fixture { + assertFailsWith { + TorrentV2DownloadSession.run(document, layout, setOf("unknown"), store, network, peerId, + buffers, state, discover = { error("Unexpected discovery") }) { + error("Unexpected session") + } + } + assertFalse(torrentFileSystem.exists(output)) + } + } +} diff --git a/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentV2EngineRateTest.kt b/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentV2EngineRateTest.kt new file mode 100644 index 000000000..dbf9ce6e9 --- /dev/null +++ b/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentV2EngineRateTest.kt @@ -0,0 +1,90 @@ +package com.linroid.ketch.torrent + +import kotlinx.coroutines.CompletableDeferred +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.async +import kotlinx.coroutines.awaitCancellation +import kotlinx.coroutines.cancelAndJoin +import kotlinx.coroutines.delay +import kotlinx.coroutines.flow.first +import kotlinx.coroutines.test.runTest +import kotlinx.coroutines.withContext +import kotlinx.coroutines.withTimeout +import okio.ByteString.Companion.toByteString +import okio.FileSystem +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertFalse +import kotlin.test.assertIs + +class TorrentV2EngineRateTest { + @Test + fun engineAndSessionLimitsGateTcpRequestsAndCanBeRemovedWhileDownloading() = runTest { + withContext(Dispatchers.Default) { + withTimeout(15_000) { + val bytes = ByteArray(32_768) { it.toByte() } + val root = sha256Digest(sha256Digest(bytes.copyOfRange(0, 16_384)) + + sha256Digest(bytes.copyOfRange(16_384, bytes.size))) + val document = TorrentV2Document.parse(Bencode.encode(mapOf("info" to mapOf( + "meta version" to 2L, "piece length" to 32_768L, + "file tree" to mapOf("a" to mapOf("" to mapOf("length" to bytes.size.toLong(), + "pieces root" to root))) + ), "piece layers" to emptyMap()))) + val output = FileSystem.SYSTEM_TEMPORARY_DIRECTORY / + "ketch-v2-rate-${InfoHash.fromBytes(torrentRandomBytes(20)).hex}" + val network = createTorrentNetwork() + val listener = network.listen(PeerEndpoint("127.0.0.1", 0)) + val peerBudget = TorrentBufferBudget(1024 * 1024) + val firstReply = CompletableDeferred() + val secondRequest = CompletableDeferred() + val server = async { + val connection = listener.accept() + try { + PeerIdentityHandshake(document.identity).respond(connection, + ByteArray(20) { 2 }.toByteString(), peerBudget) + val wire = PeerWire(connection, pieceCount = 1) + wire.send(PeerMessage.Bitfield(byteArrayOf(128.toByte()))) + wire.send(PeerMessage.Control(PeerMessage.Signal.UNCHOKE)) + assertEquals(PeerMessage.Control(PeerMessage.Signal.INTERESTED), wire.read()) + repeat(2) { index -> + val request = assertIs(wire.read()) + assertEquals(index * 16_384, request.begin) + if (index == 1) secondRequest.complete(Unit) + wire.send(PeerMessage.Piece(0, request.begin, + bytes.copyOfRange(request.begin, request.begin + request.length))) + if (index == 0) firstReply.complete(Unit) + } + } finally { connection.close() } + } + val engine = KotlinTorrentEngine(TorrentConfig(dhtEnabled = false)) + try { + engine.start() + engine.setDownloadRateLimit(1) + engine.withV2Download("rate", document, output.toString(), + discover = { it.send(listener.local); awaitCancellation() }) { session -> + session.setDownloadRateLimit(1) + session.resume() + firstReply.await() + delay(100) + assertFalse(secondRequest.isCompleted) + engine.setDownloadRateLimit(0) + delay(100) + assertFalse(secondRequest.isCompleted) + session.setDownloadRateLimit(0) + session.state.first { it == TorrentSessionState.FINISHED } + assertEquals(bytes.size.toLong(), session.verifiedBytes.value) + } + server.await() + } finally { + engine.stop() + listener.close() + network.close() + server.cancelAndJoin() + torrentFileSystem.deleteRecursively(output, mustExist = false) + } + assertEquals(0, peerBudget.allocated) + assertEquals(0, engine.admittedSessionBytes) + } + } + } +} diff --git a/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentV2SessionLoopTest.kt b/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentV2SessionLoopTest.kt index 7cfd72e54..2bcb22240 100644 --- a/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentV2SessionLoopTest.kt +++ b/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentV2SessionLoopTest.kt @@ -7,6 +7,7 @@ import kotlinx.coroutines.cancelAndJoin import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.sync.Semaphore import kotlinx.coroutines.test.StandardTestDispatcher +import kotlinx.coroutines.test.advanceTimeBy import kotlinx.coroutines.test.runCurrent import kotlinx.coroutines.test.runTest import kotlinx.coroutines.withContext @@ -129,6 +130,47 @@ class TorrentV2SessionLoopTest { } } + @Test + fun rateLimitedRequestsKeepReceivingAndHonorLiveGlobalAndTaskChanges() = runTest { + val f = Fixture(setOf("0")) + val clock = { testScheduler.currentTime } + val global = TorrentRateLimiter(1, clock) + val task = TorrentRateLimiter(1, clock) + try { + f.store.initialize() + val download = async { + PeerV2Pool.run(f.state, maxPeers = 1) { pool -> + f.attach(pool, clock) + TorrentV2CommitWorker.run(f.store, dispatcher = StandardTestDispatcher(testScheduler)) { + TorrentV2SessionLoop.download(layout, setOf("0"), f.store, pool, it, + f.buffers, f.state, maxPeers = 1, + requestDelay = { bytes, admit -> global.requestDelay(bytes, task, admit) }, + nowMs = clock) + } + } + } + try { + runCurrent() + assertEquals(1, f.connections.single().requests.size) + advanceTimeBy(1000) + runCurrent() + assertFalse(download.isCompleted) + assertEquals(1, f.connections.single().requests.size) + global.set(0) + advanceTimeBy(50) + runCurrent() + assertEquals(1, f.connections.single().requests.size) + task.set(0) + advanceTimeBy(50) + runCurrent() + download.await() + } finally { download.cancelAndJoin() } + assertTrue(f.store.completed()) + assertEquals(2, f.connections.single().requests.size) + f.checkReleased() + } finally { f.store.cleanup() } + } + @Test fun automaticallyRetriesCorruptPiecesAndWritesOnlySelectedFiles() = runTest { val selected = setOf("0")