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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
82 changes: 82 additions & 0 deletions docs/design/torrent-v2-integrity.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
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 @@ -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<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 @@ -298,23 +329,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 @@ -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 <T> withV2Download(
taskId: String,
document: TorrentV2Document,
outputPath: String,
selected: Set<String> = emptySet(),
checkpoint: TorrentV2Checkpoint? = null,
discover: suspend (SendChannel<PeerEndpoint>) -> 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<PeerEndpoint>,
Expand Down Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -24,16 +25,54 @@ 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
while (remaining > 0) {
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
Expand Down
Loading
Loading