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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 26 additions & 0 deletions docs/design/torrent-v2-integrity.md
Original file line number Diff line number Diff line change
Expand Up @@ -1092,3 +1092,29 @@ Lifecycle admission also compares the complete layout to storage's canonical hyb
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.
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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

Expand All @@ -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<RuntimeContext>
}

private val runtimeContext = RuntimeContext()
private val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default + runtimeContext)
private val shutdown = AtomicReference<Deferred<Unit>?>(null)
private val network = TorrentConnectionBudget(rawNetwork, config.maxConnections)
private val exchangeBudgets = TorrentExchangeBudgets(config)
private val budget = exchangeBudgets.transfer
Expand All @@ -55,6 +68,7 @@ internal class KotlinTorrentEngine(
private val mutex = Mutex()
private val dhtMutex = Mutex()
private val sessions = mutableMapOf<String, KotlinTorrentSession>()
private val v2Identities = mutableMapOf<String, TorrentIdentity>()
private val outputs = mutableMapOf<String, String>()
private val sessionLeases = mutableMapOf<String, TorrentBufferBudget.Lease>()
private val running = AtomicBoolean(false)
Expand All @@ -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
Expand Down Expand Up @@ -117,14 +131,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<Unit> {
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
Expand Down Expand Up @@ -294,23 +325,21 @@ 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())
.resolve(spec.outputPath).normalized()
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()))
Expand All @@ -336,6 +365,75 @@ 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.
*/
suspend fun <T> withV2Download(
taskId: String,
document: TorrentV2Document,
outputPath: String,
selected: Set<String> = emptySet(),
discover: suspend (SendChannel<PeerEndpoint>) -> Unit,
body: suspend (TorrentV2DownloadSession) -> T,
): T {
val pending = scope.async {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Prevent shutdown from joining its own v2 child

When the supplied body or discover callback calls engine.stop(), it is executing inside this engine-scope child. stop() cancels scope and then joins the scope's root job, which cannot complete until this same child returns, while the child cannot return until stop() does; the call therefore deadlocks. Run caller callbacks outside the owned child or make shutdown avoid joining the calling child.

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 8b6b1c1. Shutdown now has a single cleanup operation outside the engine job tree. An engine-owned callback requests that operation without joining its own ancestor; external stop calls await the same full cleanup barrier. Regressions call stop from both the v2 body and discovery callback, then await external stop and assert that admission returned. The combined local capability branch passes 618 JVM tests, 597 executed iOS simulator tests, the runtime guard, and both existing v1 conformance scenarios.

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)
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), 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<PeerEndpoint>,
Expand Down Expand Up @@ -451,7 +549,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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -79,3 +79,31 @@ 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,
): 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 bytes = 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" }
}
Loading
Loading