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
15 changes: 14 additions & 1 deletion docs/plans/pure-kotlin-torrent-v2-progress.md
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,8 @@ span multiple PRs; it is complete only when all of its acceptance gates have evi
| --- | --- | --- |
| 01a | `torrent-v2-01-foundation` | Pinned reference clients, executed-scenario evidence, CI gate |
| 01b | `torrent-v2-01-contracts` | Control/state/capability and compatibility contracts |
| 01c | Planned | Resource admission profiles, deterministic harness, performance baseline |
| 01c | `torrent-v2-01-budgets` | Independent metadata/transfer partitions and aggregate ceiling |
| 01d | Planned | Retained state admission, profiles, deterministic harness, baseline |

Slice 01a covers two existing v1 interoperability scenarios. Missing Transmission fails required
conformance mode; ordinary local runs omit unconfigured optional fixtures from the test plan.
Expand Down Expand Up @@ -46,3 +47,15 @@ continue their legacy behavior until the runtime adapters are implemented.
Validation: API tests pass on JVM and JavaScript; core and remote JVM implementations compile.
Mutation implementations, paginated detail endpoints, operation ledgers, and runtime adapters
remain pending; their declarations and tests ship with the respective implementation slices.

## Slice 01c

Metadata exchange now has an independent scratch reservation instead of competing with active
piece buffers. Both partitions account against a shared exchange ceiling. Configuration rejects
ceilings that cannot hold both partitions, so payload saturation cannot prevent one metadata
exchange from progressing. The existing metadata wire fixture exercises successful and failed
exchanges while the transfer partition remains fully reserved.

This is an exchange-buffer ceiling, not a total engine memory or process RSS claim. Retained
metadata/cache entries, session indexes, proof layers, disk handles, and platform allocations still
need their admission rules. Mobile/desktop production profiles and performance gates remain pending.
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,8 @@ internal class KotlinTorrentEngine(
) : TorrentEngine {
private val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
private val network = TorrentConnectionBudget(rawNetwork, config.maxConnections)
private val budget = TorrentBufferBudget(config.maxBufferedBytes)
private val exchangeBudgets = TorrentExchangeBudgets(config)
private val budget = exchangeBudgets.transfer
private val cache = TorrentMetadataCache(scope)
private val tracker = TorrentTracker(http, network)
private val peerId = torrentRandomBytes(20)
Expand Down Expand Up @@ -146,7 +147,7 @@ internal class KotlinTorrentEngine(
if (attempted.size > 4096) error("Metadata peer limit exceeded")
try {
return@coroutineScope TorrentMetadataExchange(network, config.maxMetadataBytes,
budget = budget).fetch(magnet.infoHash, endpoint)
budget = exchangeBudgets.metadata).fetch(magnet.infoHash, endpoint)
} catch (e: PrivateTorrentMagnetException) {
throw e
} catch (e: CancellationException) {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
package com.linroid.ketch.torrent

import kotlin.concurrent.atomics.AtomicBoolean
import kotlin.concurrent.atomics.AtomicInt
import kotlin.concurrent.atomics.ExperimentalAtomicApi

/** Shared across sessions; reservations precede allocating piece and connection buffers. */
@OptIn(ExperimentalAtomicApi::class)
internal class TorrentBufferBudget(
val capacity: Int,
private val parent: TorrentBufferBudget? = null,
) {
private val used = AtomicInt(0)
val allocated: Int get() = used.load()

init { require(capacity > 0) }

inner class Lease(val bytes: Int, private val parentLease: Lease?) {
private val closed = AtomicBoolean(false)
fun close() {
if (closed.compareAndSet(false, true)) {
used.fetchAndAdd(-bytes)
parentLease?.close()
}
}
}

fun reserve(bytes: Int): Lease? {
require(bytes > 0)
if (bytes > capacity) return null
val parentLease = parent?.reserve(bytes)
if (parent != null && parentLease == null) return null
while (true) {
val current = used.load()
if (bytes > capacity - current) {
parentLease?.close()
return null
}
if (used.compareAndSet(current, current + bytes)) return Lease(bytes, parentLease)
}
}
}

/** Independent partitions guarantee one metadata reservation under saturated payload load. */
internal class TorrentExchangeBudgets(config: TorrentConfig) {
private val root = TorrentBufferBudget(config.maxExchangeBytes)
val transfer = TorrentBufferBudget(config.maxBufferedBytes, root)
val metadata = TorrentBufferBudget(config.metadataExchangeBytes, root)
val allocated: Int get() = root.allocated
}
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,11 @@ data class TorrentConfig(
val maxMetadataBytes: Int = 4 * 1024 * 1024,
/** Maximum simultaneously buffered piece data across the engine. */
val maxBufferedBytes: Int = 32 * 1024 * 1024,
/**
* Combined ceiling for transfer buffers and metadata exchange scratch space. This is not a
* process RSS limit; retained metadata, session indexes, and platform overhead are separate.
*/
val maxExchangeBytes: Int = 64 * 1024 * 1024,

) {
init {
Expand All @@ -47,8 +52,15 @@ data class TorrentConfig(
require(listenPort in 0..65535) { "listenPort must be in 0..65535" }
require(maxMetadataBytes in 1..4 * 1024 * 1024)
require(maxBufferedBytes >= 16384) { "maxBufferedBytes must hold a protocol block" }
require(maxBufferedBytes.toLong() + metadataExchangeBytes <= maxExchangeBytes.toLong()) {
"maxExchangeBytes must cover transfer buffers and an independent metadata exchange"
}
}

/** One metadata exchange, including parse copies and bounded wire overhead. */
internal val metadataExchangeBytes: Int
get() = maxMetadataBytes * 4 + 256 * 1024

/** Effective policy, including compatibility with the legacy boolean. */
val effectiveUploadPolicy: TorrentUploadPolicy
get() = uploadPolicy ?: if (enableUpload) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,34 +2,6 @@ package com.linroid.ketch.torrent

import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import kotlin.concurrent.atomics.AtomicInt
import kotlin.concurrent.atomics.AtomicBoolean
import kotlin.concurrent.atomics.ExperimentalAtomicApi

/** Shared across sessions; reservations precede allocating piece and connection buffers. */
@OptIn(ExperimentalAtomicApi::class)
internal class TorrentBufferBudget(val capacity: Int) {
private val used = AtomicInt(0)
val allocated: Int get() = used.load()

init { require(capacity > 0) }

inner class Lease(val bytes: Int) {
private val closed = AtomicBoolean(false)
fun close() {
if (closed.compareAndSet(false, true)) used.fetchAndAdd(-bytes)
}
}

fun reserve(bytes: Int): Lease? {
require(bytes > 0)
while (true) {
val current = used.load()
if (bytes > capacity - current) return null
if (used.compareAndSet(current, current + bytes)) return Lease(bytes)
}
}
}

/** Serialized rarity and ownership state. No disk or network operation runs under its lock. */
internal class TorrentPieceScheduler(
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,69 @@
package com.linroid.ketch.torrent

import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFailsWith
import kotlin.test.assertNotNull
import kotlin.test.assertNull

class TorrentBufferBudgetTest {
@Test
fun childFailureRollsBackParentAndDoubleCloseDoesNotUnderflow() {
val root = TorrentBufferBudget(20)
val child = TorrentBufferBudget(10, root)
val lease = assertNotNull(child.reserve(8))
assertNull(child.reserve(3))
assertEquals(8, child.allocated)
assertEquals(8, root.allocated)
lease.close()
lease.close()
assertEquals(0, child.allocated)
assertEquals(0, root.allocated)
}

@Test
fun siblingReservationsCannotExceedAggregateCeiling() {
val root = TorrentBufferBudget(10)
val first = TorrentBufferBudget(10, root)
val second = TorrentBufferBudget(10, root)
val lease = assertNotNull(first.reserve(7))
assertNull(second.reserve(4))
assertEquals(0, second.allocated)
lease.close()
val next = assertNotNull(second.reserve(10))
assertEquals(10, root.allocated)
next.close()
assertEquals(0, root.allocated)
}

@Test
fun fullTransferPartitionCannotStarveMetadataPartition() {
val config = TorrentConfig(maxBufferedBytes = 16_384, maxMetadataBytes = 1024)
val budgets = TorrentExchangeBudgets(config)
val transfer = budgets.transfer
val metadata = budgets.metadata
val payload = assertNotNull(transfer.reserve(transfer.capacity))
val info = assertNotNull(metadata.reserve(metadata.capacity))
assertNull(transfer.reserve(1))
assertEquals(transfer.capacity + metadata.capacity, budgets.allocated)
info.close()
payload.close()
assertEquals(0, budgets.allocated)
}

@Test
fun ceilingsRejectOverflowAndInsufficientMetadataHeadroom() {
val root = TorrentBufferBudget(Int.MAX_VALUE)
val lease = assertNotNull(root.reserve(Int.MAX_VALUE))
assertNull(root.reserve(1))
lease.close()
assertEquals(0, root.allocated)
assertFailsWith<IllegalArgumentException> { root.reserve(0) }
assertFailsWith<IllegalArgumentException> {
TorrentConfig(maxBufferedBytes = Int.MAX_VALUE, maxExchangeBytes = Int.MAX_VALUE)
}
assertFailsWith<IllegalArgumentException> {
TorrentConfig(maxExchangeBytes = 32 * 1024 * 1024)
}
}
}
Original file line number Diff line number Diff line change
@@ -1,6 +1,9 @@
package com.linroid.ketch.torrent

import kotlinx.coroutines.CompletableDeferred
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.NonCancellable
import kotlinx.coroutines.awaitCancellation
import kotlinx.coroutines.async
import kotlinx.coroutines.cancelAndJoin
import kotlinx.coroutines.coroutineScope
Expand All @@ -11,10 +14,11 @@ import kotlin.test.Test
import kotlin.test.assertContentEquals
import kotlin.test.assertEquals
import kotlin.test.assertFailsWith
import kotlin.test.assertNotNull

class TorrentMetadataExchangeTest {
@Test
fun exchange_usesPerPeerIdAndAssemblesReorderedBlocks() = runTest { exchange("valid") }
fun exchange_assemblesReorderedBlocksWhileTransferBudgetIsFull() = runTest { exchange("valid") }

@Test
fun exchange_hashMismatchFailsBeforeHandoff() = runTest { exchange("corrupt") }
Expand All @@ -25,6 +29,48 @@ class TorrentMetadataExchangeTest {
@Test
fun exchange_privateMetainfoRequiresExplicitInput() = runTest { exchange("private") }

@Test
fun canceledHandshake_releasesMetadataReservationAndPreservesTransferReservation() = runTest {
withContext(Dispatchers.Default) {
withTimeout(5_000) {
val budgets = TorrentExchangeBudgets(TorrentConfig())
val payload = assertNotNull(budgets.transfer.reserve(budgets.transfer.capacity))
val network = createTorrentNetwork()
val listener = network.listen(PeerEndpoint("127.0.0.1", 0))
val ready = CompletableDeferred<Unit>()
val server = async {
val connection = listener.accept()
try {
connection.readExactly(68)
ready.complete(Unit)
awaitCancellation()
} finally { connection.close() }
}
val client = async {
TorrentMetadataExchange(network, budget = budgets.metadata)
.fetch(InfoHash.fromBytes(ByteArray(20)), listener.local)
}
try {
ready.await()
assertEquals(budgets.metadata.capacity, budgets.metadata.allocated)
client.cancelAndJoin()
assertEquals(0, budgets.metadata.allocated)
assertEquals(budgets.transfer.capacity, budgets.allocated)
} finally {
withContext(NonCancellable) {
client.cancel()
server.cancel()
network.close()
client.cancelAndJoin()
server.cancelAndJoin()
payload.close()
}
}
assertEquals(0, budgets.allocated)
}
}
}

@Test
fun extensions_updatesAreAdditiveAndCanDisableAnId() {
val extensions = PeerExtensions()
Expand All @@ -47,7 +93,9 @@ class TorrentMetadataExchangeTest {
val metadata = TorrentMetadata.fromBencode(Bencode.encode(mapOf("info" to info)))
val network = createTorrentNetwork()
val listener = network.listen(PeerEndpoint("127.0.0.1", 0))
val budget = TorrentBufferBudget(32 * 1024 * 1024)
val budgets = TorrentExchangeBudgets(TorrentConfig())
val payload = assertNotNull(budgets.transfer.reserve(budgets.transfer.capacity))
val budget = budgets.metadata
val server = async {
val connection = listener.accept()
try {
Expand Down Expand Up @@ -92,9 +140,12 @@ class TorrentMetadataExchangeTest {
}
server.await()
assertEquals(0, budget.allocated)
assertEquals(budgets.transfer.capacity, budgets.allocated)
} finally {
server.cancelAndJoin()
network.close()
payload.close()
assertEquals(0, budgets.allocated)
}
}
}
Expand Down
Loading