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
32 changes: 32 additions & 0 deletions docs/design/torrent-v2-integrity.md
Original file line number Diff line number Diff line change
Expand Up @@ -990,3 +990,35 @@ streaming recovery and large-profile memory/performance evidence remain required
Tests exercise a real engine with stale TaskStore data, corrupted payloads, same-revision counters,
replacement/binding rejection, cancellation/admission cleanup, and a failed-recovery shutdown
regression that overwrote the committed bytes before the write guard was added.

## Explicit tracker-only engine discovery

The internal engine accepts a privacy choice before resolving a magnet. Public remains the legacy
default. Tracker-only resolution requires validated supplied trackers, ignores explicit `x.pe`
endpoints, and does not start DHT discovery. It announces through one preferred tracker at a time,
tries returned metadata peers sequentially, and closes each metadata connection before a tracker
can switch. Cancellation does not fall back to public discovery. A successful tracker contact gets
a bounded best-effort stopped announce when metadata resolution ends.

Only the tracker-only path permits a hash-verified private info dictionary. Public callers still
reject private metadata, including cache hits populated by an earlier tracker-only request. Cached
info bytes do not supply endpoints or replace the current caller's tracker list. An explicit task
privacy field also disables public discovery when the resolved metadata itself is public. Peer
exchange is neither advertised nor accepted, and incoming hosts must come from tracker responses.

Real TCP tests verify tracker-authorized private metadata, failover from an unavailable tracker,
ignored explicit peers, no DHT socket attempts, private-cache rejection for public callers, missing
or invalid tracker rejection before discovery, cancellation while awaiting peers, and the task
privacy guard after public metadata resolution, including PEX rejection and incoming admission.
This is engine wiring: source/SDK selection and
persisted privacy, v2-only magnet metadata/proofs, and the broader network-policy gates remain open.
The protocol basis is [BEP 27](https://www.bittorrent.org/beps/bep_0027.html) and the info-dictionary
transfer described by [BEP 9](https://www.bittorrent.org/beps/bep_0009.html).

### Tracker-only metadata fallback review

A successful announce with no usable metadata peers now advances to the next supplied tracker
within the same resolution round. Each tracker's metadata connections close before a bounded
best-effort stopped announce and the next tracker starts. Empty responses and failed metadata
peers both have regression coverage; neither can pin resolution to the first responding tracker.
The metadata deadline and total peer-attempt bound still apply.
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ 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.delay
import kotlinx.coroutines.isActive
import kotlinx.coroutines.launch
Expand Down Expand Up @@ -146,11 +147,25 @@ internal class KotlinTorrentEngine(
}
}

override suspend fun fetchMetadata(magnetUri: String): TorrentMetadata? {
override suspend fun fetchMetadata(magnetUri: String): TorrentMetadata? =
fetchMetadata(magnetUri, TorrentDiscoveryPrivacy.PUBLIC)

override suspend fun fetchMetadata(
magnetUri: String,
privacy: TorrentDiscoveryPrivacy,
): TorrentMetadata? {
check(isRunning)
val magnet = MagnetUri.parse(magnetUri)
val restrictedTiers = if (privacy == TorrentDiscoveryPrivacy.TRACKER_ONLY) {
TrackerConfiguration.prepare(magnet.trackers.map { listOf(it) }).tiers.also {
require(it.isNotEmpty()) { "Tracker-only magnets require supplied trackers" }
}
} else emptyList()
val metadata = cache.resolve(magnet.infoHash) {
withTimeout(config.metadataTimeout) {
if (privacy == TorrentDiscoveryPrivacy.TRACKER_ONLY) {
return@withTimeout fetchTrackerOnly(magnet, restrictedTiers)
}
coroutineScope {
val peers = Channel<PeerEndpoint>(256)
val discovery = launch { discoverMagnet(magnet, peers) }
Expand Down Expand Up @@ -178,11 +193,64 @@ internal class KotlinTorrentEngine(
}
}
require(magnet.identity.matchesInfo(metadata.infoBytes)) { "Exact topic hash mismatch" }
// A cache hit must not bypass this caller's pre-discovery privacy choice.
if (metadata.isPrivate && privacy != TorrentDiscoveryPrivacy.TRACKER_ONLY) {
throw PrivateTorrentMagnetException()
}
// Cache the immutable info dictionary, retaining this caller's tracker list.
return TorrentMetadata.fromBencode(metainfoFromInfo(metadata.infoBytes,
magnet.trackers.map { listOf(it) }), config.maxMetadataBytes)
}

private suspend fun fetchTrackerOnly(
magnet: MagnetUri,
trackerTiers: List<List<String>>,
): TorrentMetadata {
val urls = trackerTiers.flatten().distinct()
var attempts = 0
while (currentCoroutineContext().isActive) {
var retrySeconds = 30L
for (url in urls) {
currentCoroutineContext().ensureActive()
// Each connection closes before switching trackers, including unsuccessful peer sets.
val tiers = TrackerTiers(listOf(listOf(url)), tracker::announce)
var contacted = false
try {
val result = attempt {
tiers.announce(TrackerAnnounce(magnet.infoHash, peerId, port, 0, 1,
event = TrackerEvent.STARTED))
} ?: continue
contacted = true
retrySeconds = maxOf(retrySeconds, result.intervalSeconds)
for (endpoint in result.peers.distinct()) {
check(++attempts <= 4096) { "Metadata peer limit exceeded" }
try {
return TorrentMetadataExchange(network, config.maxMetadataBytes,
budget = exchangeBudgets.metadata).fetch(magnet.infoHash, endpoint, trackerTiers,
TorrentDiscoveryPrivacy.TRACKER_ONLY)
} catch (error: CancellationException) {
if (!currentCoroutineContext().isActive) throw error
} catch (_: Exception) {
// Only other peers returned by these trackers are eligible retries.
}
}
} finally {
if (contacted) withContext(NonCancellable) {
withTimeoutOrNull(2000) {
attempt {
tiers.announce(TrackerAnnounce(magnet.infoHash, peerId, port, 0, 1,
event = TrackerEvent.STOPPED))
}
}
}
}
}
delay(retrySeconds.coerceAtMost(Long.MAX_VALUE / 1000) * 1000)
}
currentCoroutineContext().ensureActive()
error("Metadata resolution stopped")
}

private suspend fun discoverMagnet(magnet: MagnetUri, output: SendChannel<PeerEndpoint>) =
supervisorScope {
launch {
Expand Down Expand Up @@ -253,6 +321,7 @@ internal class KotlinTorrentEngine(
downloadThrottle = { downloadRate.acquire(it); spec.throttle(it) },
uploadThrottle = { uploadRate.acquire(it) },
trackerConfigurationBudget = exchangeBudgets.sessions,
privacy = spec.privacy,
)
sessionLeases[hash] = lease
sessions[hash] = session
Expand Down Expand Up @@ -317,7 +386,7 @@ internal class KotlinTorrentEngine(
}
}
}
if (!metadata.isPrivate) {
if (!metadata.isPrivate && spec.privacy != TorrentDiscoveryPrivacy.TRACKER_ONLY) {
spec.magnetUri?.let { uri -> launch {
MagnetUri.parse(uri).explicitPeers.forEach { text ->
resolveEndpoint(text).forEach { output.send(it) }
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ internal class KotlinTorrentSession(
private val downloadThrottle: suspend (Int) -> Unit = {},
private val uploadThrottle: suspend (Int) -> Unit = {},
private val trackerConfigurationBudget: TorrentBufferBudget? = null,
private val privacy: TorrentDiscoveryPrivacy = TorrentDiscoveryPrivacy.PUBLIC,
) : TorrentSession {
init { require(connections in 1..512) }

Expand Down Expand Up @@ -149,6 +150,8 @@ internal class KotlinTorrentSession(

fun updateTrackerStatus(status: List<TrackerStatus>) { _trackerStatus.value = status }

private val trackerRestricted = store.metadata.isPrivate ||
privacy == TorrentDiscoveryPrivacy.TRACKER_ONLY
private val privateAdmission = Mutex()
private var allowedPrivateHosts: Set<String> = emptySet()

Expand All @@ -160,7 +163,7 @@ internal class KotlinTorrentSession(
fun accept(connection: TorrentConnection): Boolean {
if (_state.value != TorrentSessionState.DOWNLOADING &&
_state.value != TorrentSessionState.SEEDING) return false
if (!store.metadata.isPrivate) return incoming.trySend(connection).isSuccess
if (!trackerRestricted) return incoming.trySend(connection).isSuccess
if (!privateAdmission.tryLock()) return false
try {
if (connection.remote.host !in allowedPrivateHosts) return false
Expand Down Expand Up @@ -242,6 +245,7 @@ internal class KotlinTorrentSession(
try {
TorrentSwarm(store, network, budget, peerId = peerId,
connections = { connectionLimit.load() }, uploadPolicy = uploadPolicy,
trackerOnly = trackerRestricted,
downloadPayload = { bytes ->
val total = received.fetchAndAdd(bytes.toLong()) + bytes
val now = clock()
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
package com.linroid.ketch.torrent

/** Chosen before discovery; metadata learned later cannot undo public lookups. */
internal enum class TorrentDiscoveryPrivacy {
PUBLIC,
TRACKER_ONLY,
}
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,12 @@ internal interface TorrentEngine {
*/
suspend fun fetchMetadata(magnetUri: String): TorrentMetadata?

/** Explicit privacy must be selected before starting magnet discovery. */
suspend fun fetchMetadata(magnetUri: String, privacy: TorrentDiscoveryPrivacy): TorrentMetadata? {
require(privacy == TorrentDiscoveryPrivacy.PUBLIC) { "Tracker-only resolution is unsupported" }
return fetchMetadata(magnetUri)
}

/**
* Adds a torrent for downloading.
*
Expand Down Expand Up @@ -86,4 +92,5 @@ internal data class TorrentTaskSpec(
val magnetUri: String? = null,
val resumeData: ByteArray? = null,
val throttle: suspend (Int) -> Unit = {},
val privacy: TorrentDiscoveryPrivacy = TorrentDiscoveryPrivacy.PUBLIC,
)
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,12 @@ internal class TorrentMetadataExchange(
hash: InfoHash,
endpoint: PeerEndpoint,
trackerTiers: List<List<String>> = emptyList(),
privacy: TorrentDiscoveryPrivacy = TorrentDiscoveryPrivacy.PUBLIC,
): TorrentMetadata = withTimeout(timeoutMs) {
require(privacy != TorrentDiscoveryPrivacy.TRACKER_ONLY ||
trackerTiers.any { it.isNotEmpty() }) {
"Tracker-only metadata requires supplied trackers"
}
var lease: TorrentBufferBudget.Lease? = null
while (lease == null) {
lease = budget.reserve(maxBytes * 4 + 256 * 1024)
Expand Down Expand Up @@ -81,7 +86,9 @@ internal class TorrentMetadataExchange(
require(InfoHash.fromBytes(sha1Digest(output)) == hash) { "Metadata hash mismatch" }
val metainfo = metainfoFromInfo(output, trackerTiers)
val metadata = TorrentMetadata.fromBencode(metainfo)
if (metadata.isPrivate) throw PrivateTorrentMagnetException()
if (metadata.isPrivate && privacy != TorrentDiscoveryPrivacy.TRACKER_ONLY) {
throw PrivateTorrentMagnetException()
}
return@withTimeout metadata
}
}
Expand Down Expand Up @@ -127,4 +134,4 @@ internal fun metainfoFromInfo(info: ByteArray, trackerTiers: List<List<String>>)
}

internal class PrivateTorrentMagnetException :
IllegalArgumentException("Private torrents require a metainfo input")
IllegalArgumentException("Private torrents require tracker-only discovery or a metainfo input")
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,9 @@ internal class TorrentSwarm(
private val onProgress: suspend (Long) -> Unit = {},
private val onCompleted: suspend () -> Unit = {},
private val allowLocalDiscovery: Boolean = false,
trackerOnly: Boolean = false,
) {
private val trackerRestricted = store.metadata.isPrivate || trackerOnly
private val uploadSlots = Semaphore(4)
private val connectedMutex = Mutex()
private val connected = mutableMapOf<Int, PeerEndpoint>()
Expand Down Expand Up @@ -72,7 +74,7 @@ internal class TorrentSwarm(
val results = Channel<Pair<PeerEndpoint, Throwable?>>(512)
val progressEvents = Channel<Unit>(Channel.CONFLATED)
val pexEvents = Channel<Pair<PeerEndpoint, PexUpdate>>(64)
val pexDirectory = TorrentPeerDirectory(store.metadata.isPrivate)
val pexDirectory = TorrentPeerDirectory(trackerRestricted)
val retryAt = mutableMapOf<PeerEndpoint, Long>()
val now = monotonicClock()
var discoveryClosed = false
Expand Down Expand Up @@ -224,7 +226,7 @@ internal class TorrentSwarm(
var metadataServed = 0
var metadataWindow = TimeSource.Monotonic.markNow()
if (handshake.extensions) {
wire.send(PeerExtensions.handshake(store.metadata, pex = !store.metadata.isPrivate))
wire.send(PeerExtensions.handshake(store.metadata, pex = !trackerRestricted))
}
val state = PeerProtocolState(store.pieceCount, maxPending = 16)
val advertised = store.verifiedPieces()
Expand Down Expand Up @@ -337,8 +339,8 @@ internal class TorrentSwarm(
}
}
} else if (message.id == PeerExtensions.PEX) {
require(!store.metadata.isPrivate) {
"Private torrent peer exchange is forbidden"
require(!trackerRestricted) {
"Tracker-restricted peer exchange is forbidden"
}
val update = exchange.receive(message.payload)
val added = update.added.filter { peer ->
Expand All @@ -358,7 +360,7 @@ internal class TorrentSwarm(
uploadSlot = true
wire.send(PeerMessage.Control(PeerMessage.Signal.UNCHOKE))
}
if (!store.metadata.isPrivate && extensions.id("ut_pex") != 0 && exchange.due()) {
if (!trackerRestricted && extensions.id("ut_pex") != 0 && exchange.due()) {
val contacts = connectedMutex.withLock { connected.values.toSet() }
.filter { it != connection.remote && (allowLocalDiscovery ||
numericAddress(it.host)?.let(::publicTorrentAddress) == true) }.toSet()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -40,8 +40,13 @@ class SessionTrackerEditTest {

private class Provider : ForwardingFileSystem(torrentFileSystem) {
var failMove = false
var failNextMove = false
var afterMove: (() -> Unit)? = null
override fun atomicMove(source: Path, target: Path) {
if (target.name == "checkpoint" && failNextMove) {
failNextMove = false
throw IOException("Injected one-time rename failure")
}
if (target.name == "checkpoint" && failMove) throw IOException("Injected rename failure")
super.atomicMove(source, target)
if (target.name == "checkpoint") afterMove?.invoke()
Expand Down Expand Up @@ -165,14 +170,14 @@ class SessionTrackerEditTest {
session.resume()
discoveries.first { it.size == 1 }
session.state.first { it == TorrentSessionState.SEEDING }
provider.failMove = true
// Fail only the edit; the restarted session must be able to checkpoint again.
provider.failNextMove = true
assertFailsWith<IOException> { session.replaceTrackers(next) }
discoveries.first { it.size == 2 }
assertEquals(listOf(old, old), discoveries.value)
assertEquals(retained, state.allocated)
assertEquals(old, session.trackerTiers())
assertEquals(1L, session.trackerConfiguration().revision)
provider.failMove = false
assertTrue(session.replaceTrackers(next))
discoveries.first { it.size == 3 }
assertEquals(next, discoveries.value.last())
Expand Down
Loading
Loading