diff --git a/docs/plans/pure-kotlin-torrent-v2-progress.md b/docs/plans/pure-kotlin-torrent-v2-progress.md index 974b83584..a9e5f7ec0 100644 --- a/docs/plans/pure-kotlin-torrent-v2-progress.md +++ b/docs/plans/pure-kotlin-torrent-v2-progress.md @@ -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. @@ -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. 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 34d027421..95acb6027 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 @@ -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) @@ -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) { diff --git a/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentBufferBudget.kt b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentBufferBudget.kt new file mode 100644 index 000000000..e254e296a --- /dev/null +++ b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentBufferBudget.kt @@ -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 +} diff --git a/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentConfig.kt b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentConfig.kt index b2ebae36a..910502d7e 100644 --- a/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentConfig.kt +++ b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentConfig.kt @@ -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 { @@ -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) { diff --git a/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentPieceScheduler.kt b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentPieceScheduler.kt index a627bd311..de59e5144 100644 --- a/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentPieceScheduler.kt +++ b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentPieceScheduler.kt @@ -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( diff --git a/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentBufferBudgetTest.kt b/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentBufferBudgetTest.kt new file mode 100644 index 000000000..c42e9e2fa --- /dev/null +++ b/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentBufferBudgetTest.kt @@ -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 { root.reserve(0) } + assertFailsWith { + TorrentConfig(maxBufferedBytes = Int.MAX_VALUE, maxExchangeBytes = Int.MAX_VALUE) + } + assertFailsWith { + TorrentConfig(maxExchangeBytes = 32 * 1024 * 1024) + } + } +} diff --git a/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentMetadataExchangeTest.kt b/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentMetadataExchangeTest.kt index ae34b9392..2315dabfb 100644 --- a/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentMetadataExchangeTest.kt +++ b/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentMetadataExchangeTest.kt @@ -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 @@ -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") } @@ -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() + 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() @@ -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 { @@ -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) } } }