diff --git a/.github/workflows/tests.yml b/.github/workflows/tests.yml index a96b3a56a..ea6184127 100644 --- a/.github/workflows/tests.yml +++ b/.github/workflows/tests.yml @@ -4,7 +4,7 @@ on: push: branches: [ main ] pull_request: - branches: [ main, "codex/torrent-*" ] + branches: [ main, "codex/torrent-*", "torrent-v2-*" ] concurrency: group: ${{ github.workflow }}-${{ github.ref }} @@ -18,7 +18,7 @@ permissions: jobs: jvm-tests: name: JVM Tests - runs-on: ubuntu-latest + runs-on: ubuntu-24.04 steps: - uses: actions/checkout@v6 @@ -29,13 +29,41 @@ jobs: - uses: gradle/actions/setup-gradle@v5 - - name: Install independent torrent client - run: sudo apt-get update && sudo apt-get install -y transmission-daemon + - name: Install fixture build prerequisites + run: | + sudo apt-get update + sudo apt-get install -y cmake build-essential libcurl4-openssl-dev libssl-dev + + - name: Cache pinned Transmission fixture + uses: actions/cache@v4 + with: + path: ${{ runner.temp }}/ketch-transmission + key: transmission-${{ runner.os }}-${{ runner.arch }}-${{ hashFiles('test-fixtures/torrent/clients.properties', 'tools/torrent/build_transmission.py') }} + + - name: Build authenticated independent torrent client + run: python3 tools/torrent/build_transmission.py --output "$RUNNER_TEMP/ketch-transmission" + + - name: Verify conformance gate behavior + run: python3 -m unittest discover -s tools/torrent -p 'test_*.py' - name: Run JVM tests env: - TRANSMISSION_DAEMON: /usr/bin/transmission-daemon - run: ./gradlew jvmTest :library:torrent:verifyNoNativeTorrentRuntime + TRANSMISSION_DAEMON: ${{ runner.temp }}/ketch-transmission/build/daemon/transmission-daemon + run: ./gradlew jvmTest :library:torrent:verifyNoNativeTorrentRuntime -PtorrentConformance=true + + - name: Require executed interoperability scenarios + run: | + python3 tools/torrent/verify_conformance.py \ + --reports library/torrent/build/test-results/jvmTest \ + --revision "$GITHUB_SHA" \ + --output library/torrent/build/reports/conformance/executed.json + + - name: Upload conformance evidence + uses: actions/upload-artifact@v7 + with: + name: torrent-conformance + path: library/torrent/build/reports/conformance/executed.json + if-no-files-found: error - name: Upload test results if: always() @@ -173,6 +201,24 @@ jobs: name: test-results-torrent-device path: library/torrent/build/outputs/androidTest-results/**/TEST-*.xml + required-checks: + name: Required Checks + runs-on: ubuntu-24.04 + needs: [ jvm-tests, android-tests, ios-tests, js-tests, torrent-desktop-tests, torrent-device-tests, publish-results ] + if: always() + steps: + - name: Require every test job to succeed on this revision + env: + JOB_RESULTS: ${{ toJSON(needs) }} + run: | + python3 - <<'PY' + import json, os + results = json.loads(os.environ['JOB_RESULTS']) + failed = {name: value['result'] for name, value in results.items() + if value['result'] != 'success'} + assert results and not failed, f'Required checks did not pass: {failed}' + PY + publish-results: name: Publish Test Results runs-on: ubuntu-latest @@ -183,9 +229,12 @@ jobs: uses: actions/download-artifact@v8 with: pattern: test-results-* - merge-multiple: true + # Host jobs use identical report paths. Keep each artifact in its own directory. + merge-multiple: false - name: Publish test results uses: EnricoMi/publish-unit-test-result-action@v2 with: files: '**/TEST-*.xml' + action_fail: true + action_fail_on_inconclusive: true diff --git a/docs/design/torrent-control-contract.md b/docs/design/torrent-control-contract.md new file mode 100644 index 000000000..afb508516 --- /dev/null +++ b/docs/design/torrent-control-contract.md @@ -0,0 +1,102 @@ +# Torrent control protocol, version 1 + +This contract implements the API boundary decision in [roadmap #162][roadmap]. The types live in +`library:api`, so WebAssembly and remote clients can use them without linking a torrent engine. +This PR defines inspection, capability negotiation, ordering, and command preconditions. Runtime +wiring and the typed mutation methods ship with their capabilities, with SDK/HTTP/SSE parity checked +in roadmap step 27. A contract declaration is not an implemented runtime capability. + +## Availability and negotiation + +`KetchApi.torrents` is nullable and defaults to null for existing backend implementations. Null means +no typed control protocol, not no torrent downloader. Local and remote adapters expose a controller +only when they implement inspection. They advertise only executable features. The backend's +platform determines capabilities: a browser connected to a daemon can have capabilities unavailable +in a local browser. Protocol major versions other than 1 disable all known features in this SDK. +Unknown capability strings survive decoding. New mandatory semantics require a new major version; +optional fields may be added without changing the interpretation of existing fields. + +Page sizes and subscription limits are explicit bounded ceilings. Oversized requests fail rather +than being silently truncated. Cursors bind to task, authenticated principal, filters, sort, and +revision. Mutations that invalidate a page return a stale-cursor error; no mixed-revision file list. +Peer addresses and other sensitive details are requested separately, not broadcast in every event. + +## State and completion + +`TorrentSnapshot` is a bounded aggregate. Metadata-unavailable counters are null, not invented zeroes. +Selected verified bytes, wanted bytes, payload received, uploaded payload, discarded payload, and +protocol overhead are separate. Virtual padding never counts toward user payload. A recheck can +reduce verified progress. Counters cannot exceed their content bounds or become negative. + +A selection has a monotonically increasing generation. Completion belongs to that generation and +can coexist with seeding. Expanding a completed selection requires an explicit restart; existing +`awaitCompletion` callers retain their original generation. Empty selections can complete once +metadata is known. They do not imply possession of the torrent's payload. Download queue slots and +seed slots remain separate. Legacy upload remains disabled unless explicitly enabled; the named +production profile can enable uploads while downloading. Post-completion seeding remains opt-in. + +## Revision and reconnect ordering + +Every published task state has an opaque catalog epoch and a nonnegative sequence. Sequences never +wrap. A backend changes the epoch if it loses monotonic state, including restoration of an older +catalog. Epoch strings are compared for equality, never sorted. Full snapshots can jump forward; +deltas require their exact predecessor. Older or duplicate same-epoch messages are ignored. Gaps +in delta history and epoch changes require an authoritative resync. + +`compareIncoming` implements this merge decision. It assumes messages belong to the active +connection generation. Adapters separately reject late responses from old connections. Reconnect +cancels the old subscription, fetches authoritative state, and installs the new epoch under a new +connection generation. An old HTTP response cannot overwrite a newer event in the same epoch. +An HTTP response from an obsolete connection cannot reset the new epoch. + +Inspection streams emit full snapshots and conflate slow consumers to the latest snapshot. Null +means a tombstone and completes the task stream. A tombstone is terminal for that task ID, which is +never reused. The client rejects all subsequent responses for that removed task in the current +connection generation. Transport failures terminate observation; reconnect is explicit. Persisted +idempotency outcomes for removed tasks remain available during the supported retry window. + +## Mutations and operations + +Each existing-task mutation carries `TorrentCommandContext`: an idempotency key and the expected +revision. Scope the key to principal and task. Check the persisted retry ledger before revision +validation: an exact retry returns its original result, even if state has since advanced. Reusing +a key for a different canonical request fails. New requests with stale revisions return a conflict +and current revision, without partial mutation. Concurrent duplicates execute only once. + +Long-running operations return an operation ID with queued/running/succeeded/failed/canceled state, +progress, and cancellation support. Acceptance never implies completion. A mutation's committed +outcome and idempotency record are persisted atomically. Ledger capacity is bounded; reject new +admissions rather than evicting entries still inside the advertised retry window. After expiration, +a retry fails explicitly instead of unexpectedly executing an old destructive command again. + +The typed command families and required semantics are: + +| Command family | Contract | +| --- | --- | +| Selection | Stable file IDs; priorities; sequential mode; verified-range deadlines; explicit restart | +| Transfer | Download/upload limits; pause/resume; runtime admission applies to both directions | +| Seeding | Ratio and duration goals; start/stop; separate seed queue; no implicit upload opt-in | +| Integrity | Recheck/import as cancelable operations; publish only verified state | +| Trackers | Edit ordered tiers and reannounce; retain private/proxy policy and credential context | +| Storage | Move/rename as recoverable operations; reject unsafe paths and ownership conflicts | +| Removal | Explicit keep-data or remove-owned-data policy; never delete unrelated files | +| Export | Export metainfo/magnets without publication, tracker requests, or starting seeding | +| Creation | Stable source snapshot, format and piece policy; cancelable; no implicit publication | +| Streaming | Authenticated task/file/range reads; only verified ranges; bounded lifetime and buffers | + +Creation has no existing task revision; its request is scoped to the principal's creation ledger. +Operation cancellation has its own idempotency key and operation revision. Errors distinguish +unsupported capability, invalid input, conflict, policy denial, resource exhaustion, storage failure, +integrity failure, not found, and cancellation. They must not expose tracker credentials or raw +untrusted paths. Unsupported requests fail before any side effect or network access. + +## Migration gates + +Existing `KetchApi` implementations compile with the default null controller. New clients connected +to old daemons keep legacy downloads working and hide unsupported controls. Adding this API does not +change the existing resume format. Checkpoint v2 migration retains old state until verified commit, +rehashes legacy v1 data, and preserves unsupported native blobs for explicit recovery. No optimistic +conversion may claim unverified bytes. Format identity and checkpoint implementation remain in +roadmap steps 04–09. These gates are pending until exercised by their implementation tests. + +[roadmap]: https://github.com/linroid/Ketch/issues/162 diff --git a/docs/design/torrent-resource-accounting.md b/docs/design/torrent-resource-accounting.md new file mode 100644 index 000000000..5e092f17a --- /dev/null +++ b/docs/design/torrent-resource-accounting.md @@ -0,0 +1,38 @@ +# Torrent resource accounting + +The v2 roadmap requires simultaneous limits for memory, peers, active tasks, payload handles, +metadata, files, pieces, and hash layers. The current foundation provides some of these limits; +it does not yet implement or certify the complete mobile and desktop production profiles. + +## Payload handles + +`TorrentConfig.maxOpenPayloadFiles` bounds aggregate open payload handles across an engine's +sessions. It defaults to 32 and accepts values from 1 through 128. The engine gives every piece +store the same coroutine semaphore. Waiting is cancellable and happens before dispatch to the +blocking I/O executor, so a waiting session does not occupy an I/O thread or open a payload file. + +A piece-store operation opens at most one payload handle at a time, including the boundary-piece +sidecars used for partial selection. A semaphore permit spans the whole storage operation, rather +than an individual open/close pair. Initialization, verification, reads, commits, final truncation, +and cleanup consequently share admission. Long scans can hold a permit across multiple sequential +file opens; the implementation favors a strict bound over maximum concurrency. + +Cancellation cannot release a permit while a blocking provider call is still executing. The +permit returns only after the I/O block returns or throws and its scoped handles close. Canceling +a waiter removes that wait without taking a permit. A provider failure returns the permit so +another torrent can proceed. This does not make a blocked filesystem provider cancellable or prove +the production pause/stop latency gate. + +Ownership journals and checkpoints can open an additional non-payload handle while a payload +handle is open. They are outside this payload ceiling, as are sockets and platform/runtime file +descriptors. Process descriptor counts must be measured separately. Recovery-only stores used by +source cleanup read ownership records and remove files; they do not open payload handles. + +## Remaining production work + +The current buffer, metadata-cache, and session admission estimates are not a complete accounting +of engine working memory and cannot establish a process RSS ceiling. Raw metainfo parsing still +has limits below the proposed desktop profile. Separate seed admission, external hash-layer +storage, complete parser/index accounting, and the named production profiles remain roadmap work. +Physical-device, lifecycle, memory, descriptor, and soak evidence must cover the completed runtime +before either profile is described as production-ready. diff --git a/docs/plans/pure-kotlin-torrent-v2-progress.md b/docs/plans/pure-kotlin-torrent-v2-progress.md new file mode 100644 index 000000000..b1a734174 --- /dev/null +++ b/docs/plans/pure-kotlin-torrent-v2-progress.md @@ -0,0 +1,101 @@ +# Pure Kotlin torrent v2 implementation progress + +The approved scope and acceptance gates are tracked in +[roadmap #162](https://github.com/linroid/Ketch/issues/162). +Implementation uses focused stacked PRs, each based on the preceding branch. A roadmap row may +span multiple PRs; it is complete only when all of its acceptance gates have evidence. + +## Foundation stack + +| Slice | Branch | Scope | +| --- | --- | --- | +| 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 | `torrent-v2-01-budgets` | Independent metadata/transfer partitions and aggregate ceiling | +| 01d | `torrent-v2-01-admission` | Retained metadata cache admission and lifecycle cleanup | +| 01e | `torrent-v2-01-session-admission` | Session admission before storage/index construction | +| Remaining 01 | Planned | Production profiles, deterministic harness, performance 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. +The JVM CI lane verifies actual executions and publishes evidence tied to its tested revision. +`Required Checks` aggregates the existing platform lanes; repository branch protection must select +that check separately. Adding a job does not configure GitHub branch protection. + +Roadmap step 01 remains incomplete. This harness is not evidence of v2/hybrid compatibility, +production performance, physical-device behavior, or completion of any later roadmap step. +Steps 02–30 remain pending. Runtime behavior and the legacy upload default are unchanged by 01a. + +## Validation environment for 01a + +- Baseline: merged v1 tree at `06e8e2cb`. +- Local host: macOS 26.6.2 (25G83), arm64; OpenJDK 21.0.11; Gradle 9.7.1. +- Reference clients: source-built Transmission 4.1.3 and JVM test-only libtorrent4j 2.1.0-39. +- Fixture source digest and scenario identities: `test-fixtures/torrent/`. +- The existing v1 JVM suite and native runtime dependency guard pass locally. +- Six verifier tests cover missing, skipped, failed, duplicate, and executed report behavior. + +CI artifacts provide revision-specific results. Local fixture timings are correctness test durations, +not throughput or memory benchmarks; the production performance baseline is still pending. + +## Slice 01b + +The [control protocol contract](../design/torrent-control-contract.md) specifies capabilities, +revision/reconnect ordering, completion generations, mutation semantics, and compatibility gates. +The API module supplies bounded inspection models, capability negotiation, command preconditions, +and revision merge decisions. `KetchApi.torrents` defaults to null; existing local/remote backends +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. + +## Slice 01d + +Cached metadata now holds a budget lease until eviction, replacement, or shutdown. Conservative +weights include retained arrays, file records, strings, and a container allowance. Results that +cannot be admitted are returned to the caller without being retained. Cache capacity is a separate +partition within the same ceiling, preserving scratch headroom for metadata exchange even when +cache and transfer partitions are full. + +Explicit shutdown rejects new cache work, cancels shared fetches, and awaits their finalizers +outside the cache mutex. Owner cancellation also clears retained entries. Tests cover replacement, +least-recently-used eviction, parent pressure, oversized entries, repeated close, pending fetch +cancellation, and owner teardown. Caller-owned metadata, session state, and process overhead remain +outside this cache-retention accounting; production profile and total-memory gates remain pending. + +## Slice 01e + +Session admission applies file count, piece count, per-session estimated size, and aggregate state +capacity before constructing storage or decoding a checkpoint. The allowance covers retained +metadata, file/path indexes, scheduler arrays/candidate lists, checkpoint parsing, and checking +scratch buffers. Existing peer/wire reservations remain in the transfer partition. + +Reservations survive pause and failed deletion, return on construction failure, and return after +successful detachment on removal or after the entire runtime's jobs finish during shutdown. Tests verify rejection before filesystem creation, +aggregate exhaustion, failed-construction rollback, pause retention, removal/readmission, and +non-suspending shutdown. The combined ceiling includes session admission while preserving metadata +headroom even if every other partition is full. + +These are conservative admission allowances, not measured total heap or RSS bounds. New format +limits, runtime profiles, platform overhead, deterministic transport fault coverage, and baseline +performance evidence still require the remaining roadmap work. No production release gate is +checked off from these admission tests alone. + +Review follow-up for 01e: closing a session's network jobs does not end runtime ownership when +file cleanup fails. A runtime ledger now retains admission until successful map detachment. +Failed deletion remains retryable instead of being treated as already completed on the next call. +The regression test proves two failed cleanup attempts remain charged and block new admission, +while explicit keep-data removal returns credit without touching storage. diff --git a/library/api/src/commonMain/kotlin/com/linroid/ketch/api/KetchApi.kt b/library/api/src/commonMain/kotlin/com/linroid/ketch/api/KetchApi.kt index 8ef100736..8f138eefb 100644 --- a/library/api/src/commonMain/kotlin/com/linroid/ketch/api/KetchApi.kt +++ b/library/api/src/commonMain/kotlin/com/linroid/ketch/api/KetchApi.kt @@ -1,5 +1,6 @@ package com.linroid.ketch.api +import com.linroid.ketch.api.torrent.TorrentController import kotlinx.coroutines.flow.StateFlow /** @@ -11,6 +12,12 @@ interface KetchApi { /** Human-readable label: "Core" or "Remote · host:port". */ val backendLabel: String + /** + * Optional typed torrent controls exposed by this backend. Null preserves compatibility with + * backends that support legacy torrent downloads but have not implemented the control protocol. + */ + val torrents: TorrentController? get() = null + /** Reactive task list updated on any state change. */ val tasks: StateFlow> diff --git a/library/api/src/commonMain/kotlin/com/linroid/ketch/api/torrent/TorrentCapabilities.kt b/library/api/src/commonMain/kotlin/com/linroid/ketch/api/torrent/TorrentCapabilities.kt new file mode 100644 index 000000000..24d2d69cd --- /dev/null +++ b/library/api/src/commonMain/kotlin/com/linroid/ketch/api/torrent/TorrentCapabilities.kt @@ -0,0 +1,54 @@ +package com.linroid.ketch.api.torrent + +import kotlinx.serialization.Serializable + +/** Stable capability names; unknown advertised names are retained for forward compatibility. */ +enum class TorrentCapability(val wireName: String) { + INSPECT("inspect"), + FILE_SELECTION("file-selection"), + STREAMING("verified-streaming"), + TRANSFER_LIMITS("transfer-limits"), + SEEDING("seeding"), + RECHECK("recheck"), + TRACKERS("trackers"), + RELOCATE("relocate"), + RENAME("rename"), + IMPORT("import"), + EXPORT("export"), + CREATE("create"), + REMOVE_OWNED_DATA("remove-owned-data"), + V1("v1"), + V2("v2"), + HYBRID("hybrid"), + UTP("utp"), + ENCRYPTION("mse-pe"), + PROXY("proxy"), +} + +/** + * Capabilities of the connected backend, not the frontend platform. + * + * Names are strings so a newer backend can advertise features an older SDK does not know. + * An incompatible major version must never be interpreted as supporting a known command. + * Limits are negotiated ceilings, not a promise that admission will succeed under current load. + * [backgroundTransfers] describes unattended execution; foreground transfers may still work. + */ +@Serializable +data class TorrentCapabilities( + val protocolMajor: Int = 1, + val names: Set = emptySet(), + val maxPageSize: Int = 100, + val maxSubscriptions: Int = 16, + val backgroundTransfers: Boolean = false, +) { + init { + require(protocolMajor > 0) + require(names.size <= 128 && names.all { it.length in 1..64 }) + require(maxPageSize in 1..1000) + require(maxSubscriptions in 1..1024) + } + + /** Whether this SDK can safely use the advertised feature. Unknown major versions fail closed. */ + fun supports(capability: TorrentCapability): Boolean = + protocolMajor == 1 && capability.wireName in names +} diff --git a/library/api/src/commonMain/kotlin/com/linroid/ketch/api/torrent/TorrentController.kt b/library/api/src/commonMain/kotlin/com/linroid/ketch/api/torrent/TorrentController.kt new file mode 100644 index 000000000..417cd5f59 --- /dev/null +++ b/library/api/src/commonMain/kotlin/com/linroid/ketch/api/torrent/TorrentController.kt @@ -0,0 +1,34 @@ +package com.linroid.ketch.api.torrent + +import kotlinx.coroutines.flow.Flow + +/** + * Backend-owned torrent inspection and capability negotiation. + * + * Obtain this optional controller from KetchApi. A null controller means the backend does not + * expose this protocol; it does not imply that legacy torrent downloads are unavailable. + * Runtime implementations and mutation commands are introduced with their roadmap capabilities. + * Implementations advertise only features they can execute, never planned features. + */ +interface TorrentController { + /** Negotiate before subscribing or issuing feature-specific requests. */ + suspend fun capabilities(): TorrentCapabilities + + /** + * Get an authoritative summary, or null if the task no longer exists or is not a torrent. + * Callers must merge concurrent HTTP/event results using [TorrentRevision.compareIncoming]. + */ + suspend fun snapshot(taskId: String): TorrentSnapshot? + + /** + * Observe full snapshots for one task, starting with its current value. Null is a tombstone and + * completes the stream. Unknown tasks emit null and complete. The backend bounds subscriptions; + * each slow subscriber receives the latest full snapshot, not an unbounded queue of updates. + * + * The adapter cancels old subscriptions before reconnecting, obtains a new authoritative + * snapshot, + * and discards all late responses from the old connection generation. Revisions within an epoch + * never decrease. Transport errors terminate the stream and require explicit reconnect. + */ + fun observe(taskId: String): Flow +} diff --git a/library/api/src/commonMain/kotlin/com/linroid/ketch/api/torrent/TorrentRevision.kt b/library/api/src/commonMain/kotlin/com/linroid/ketch/api/torrent/TorrentRevision.kt new file mode 100644 index 000000000..d0913300f --- /dev/null +++ b/library/api/src/commonMain/kotlin/com/linroid/ketch/api/torrent/TorrentRevision.kt @@ -0,0 +1,59 @@ +package com.linroid.ketch.api.torrent + +import kotlinx.serialization.Serializable + +/** + * Revision within one backend state incarnation. The opaque [epoch] changes whenever the backend + * cannot preserve monotonic revisions (including restoring an older catalog). Counters never wrap. + * Epochs cannot be ordered; a different epoch requires an authoritative reconnect snapshot. + */ +@Serializable +data class TorrentRevision(val epoch: String, val sequence: Long) { + init { + require(epoch.isNotBlank() && epoch.length <= 128) + require(sequence >= 0) + } +} + +/** Result of comparing an incoming snapshot/event against the current task view. */ +enum class TorrentRevisionDecision { + APPLY, + IGNORE, + RESYNC, +} + +/** + * Decide whether to apply a revision from the current connection generation. + * + * A delta requires an exact predecessor; full snapshots may jump forward. Duplicate/older data + * never replaces newer state. A changed epoch requires reconnect initialization, not implicit + * replacement. The caller must separately discard responses from obsolete connection generations. + */ +fun TorrentRevision.compareIncoming( + incoming: TorrentRevision, + predecessor: TorrentRevision? = null, +): TorrentRevisionDecision = when { + incoming.epoch != epoch -> TorrentRevisionDecision.RESYNC + incoming.sequence <= sequence -> TorrentRevisionDecision.IGNORE + predecessor != null && predecessor != this -> TorrentRevisionDecision.RESYNC + else -> TorrentRevisionDecision.APPLY +} + +/** + * Optimistic concurrency and retry identity for one mutation of an existing task. + * + * The backend scopes [idempotencyKey] to the authenticated principal and task. An exact retry + * returns + * the original outcome; reuse with a different command fails. A new command requires the current + * [expectedRevision]. Check the retry ledger before checking the revision. Long-running operations + * return their operation ID; accepting a command does not imply completion. + */ +@Serializable +data class TorrentCommandContext( + val idempotencyKey: String, + val expectedRevision: TorrentRevision, +) { + init { + require(idempotencyKey.isNotBlank() && idempotencyKey.length <= 128) + } +} diff --git a/library/api/src/commonMain/kotlin/com/linroid/ketch/api/torrent/TorrentState.kt b/library/api/src/commonMain/kotlin/com/linroid/ketch/api/torrent/TorrentState.kt new file mode 100644 index 000000000..fe11b77a8 --- /dev/null +++ b/library/api/src/commonMain/kotlin/com/linroid/ketch/api/torrent/TorrentState.kt @@ -0,0 +1,76 @@ +package com.linroid.ketch.api.torrent + +import kotlinx.serialization.Serializable + +/** Transfer lifecycle is independent of the download task's selected-data completion. */ +@Serializable +enum class TorrentActivity { + RESOLVING, + CHECKING, + QUEUED, + DOWNLOADING, + SEEDING, + PAUSED, + STOPPED, + FAILED, +} + +/** + * Aggregate counters. Payload excludes padding and protocol overhead; retries/discards are + * separate. + * [selectedVerifiedBytes] can decrease after selection changes or a failed recheck. Completion is + * scoped to the snapshot's selection generation. Ratio is undefined when downloaded payload is + * zero. + */ +@Serializable +data class TorrentCounters( + val totalPayloadBytes: Long, + val wantedBytes: Long, + val selectedVerifiedBytes: Long, + val receivedPayloadBytes: Long, + val uploadedPayloadBytes: Long, + val discardedPayloadBytes: Long, + val protocolBytes: Long, + val downloadBytesPerSecond: Long, + val uploadBytesPerSecond: Long, + val seedSeconds: Long, +) { + init { + require(totalPayloadBytes >= 0 && wantedBytes in 0..totalPayloadBytes) + require(selectedVerifiedBytes in 0..wantedBytes) + require(receivedPayloadBytes >= 0 && uploadedPayloadBytes >= 0) + require(discardedPayloadBytes >= 0 && protocolBytes >= 0) + require(downloadBytesPerSecond >= 0 && uploadBytesPerSecond >= 0 && seedSeconds >= 0) + } +} + +/** + * Bounded task summary. Peer/file lists are paginated separately at the same revision. + * Metadata may be absent while resolving; [counters] then remains null rather than implying zero. + * A completed selection does not imply full-content availability or stopped seeding. + */ +@Serializable +data class TorrentSnapshot( + val taskId: String, + val revision: TorrentRevision, + val activity: TorrentActivity, + val selectionGeneration: Long, + val selectionComplete: Boolean, + val counters: TorrentCounters? = null, +) { + init { + require(taskId.isNotBlank() && taskId.length <= 128) + require(selectionGeneration >= 0) + require(!selectionComplete || counters != null && + counters.selectedVerifiedBytes == counters.wantedBytes) + } +} + +/** Bounded page request with opaque cursors, invalidated by relevant mutations. */ +@Serializable +data class TorrentPageRequest(val limit: Int = 100, val cursor: String? = null) { + init { + require(limit in 1..1000) + require(cursor == null || cursor.isNotBlank() && cursor.length <= 4096) + } +} diff --git a/library/api/src/commonTest/kotlin/com/linroid/ketch/api/torrent/TorrentContractsTest.kt b/library/api/src/commonTest/kotlin/com/linroid/ketch/api/torrent/TorrentContractsTest.kt new file mode 100644 index 000000000..13f7930b2 --- /dev/null +++ b/library/api/src/commonTest/kotlin/com/linroid/ketch/api/torrent/TorrentContractsTest.kt @@ -0,0 +1,73 @@ +package com.linroid.ketch.api.torrent + +import kotlinx.serialization.json.Json +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertFailsWith +import kotlin.test.assertFalse +import kotlin.test.assertTrue + +class TorrentContractsTest { + @Test + fun capabilities_unknownNamesSurviveButUnknownMajorDoesNotEnableCommands() { + val capabilities = Json.decodeFromString( + """{"names":["v2","future-feature"]}""" + ) + assertTrue(capabilities.supports(TorrentCapability.V2)) + assertFalse(capabilities.supports(TorrentCapability.HYBRID)) + assertTrue("future-feature" in capabilities.names) + assertFalse(capabilities.copy(protocolMajor = 2).supports(TorrentCapability.V2)) + } + + @Test + fun snapshotOrdering_oldHttpResponseCannotReplaceNewerEvent() { + val current = TorrentRevision("catalog-a", 10) + assertEquals(TorrentRevisionDecision.IGNORE, + current.compareIncoming(TorrentRevision("catalog-a", 9))) + assertEquals(TorrentRevisionDecision.IGNORE, current.compareIncoming(current)) + assertEquals(TorrentRevisionDecision.APPLY, + current.compareIncoming(TorrentRevision("catalog-a", 12))) + } + + @Test + fun deltaOrdering_gapsAndBackendRestartsRequireResync() { + val current = TorrentRevision("catalog-a", 10) + val next = TorrentRevision("catalog-a", 12) + assertEquals(TorrentRevisionDecision.APPLY, current.compareIncoming(next, current)) + assertEquals(TorrentRevisionDecision.RESYNC, + current.compareIncoming(next, TorrentRevision("catalog-a", 11))) + assertEquals(TorrentRevisionDecision.RESYNC, + current.compareIncoming(TorrentRevision("catalog-b", 1))) + assertEquals(TorrentRevisionDecision.RESYNC, + current.compareIncoming(next, TorrentRevision("catalog-b", 10))) + } + + @Test + fun malformedWireValuesAreRejected() { + assertFailsWith { + Json.decodeFromString("""{"epoch":"a","sequence":-1}""") + } + assertFailsWith { TorrentPageRequest(limit = 1001) } + assertFailsWith { TorrentPageRequest(cursor = "") } + assertFailsWith { + TorrentCommandContext(" ", TorrentRevision("a", 0)) + } + assertFailsWith { + TorrentCapabilities(names = setOf("x".repeat(65))) + } + } + + @Test + fun completionRequiresVerifiedSelectionButDoesNotRequireWholeTorrent() { + val counters = TorrentCounters(100, 40, 40, 50, 0, 10, 15, 0, 0, 0) + val snapshot = TorrentSnapshot("task", TorrentRevision("a", 1), TorrentActivity.SEEDING, + 2, true, counters) + assertTrue(snapshot.selectionComplete) + assertFailsWith { + snapshot.copy(counters = counters.copy(selectedVerifiedBytes = 39)) + } + assertFailsWith { snapshot.copy(counters = null) } + assertFailsWith { counters.copy(wantedBytes = 101) } + assertFailsWith { counters.copy(selectedVerifiedBytes = -1) } + } +} diff --git a/library/torrent/build.gradle.kts b/library/torrent/build.gradle.kts index 157c07abe..d76b57270 100644 --- a/library/torrent/build.gradle.kts +++ b/library/torrent/build.gradle.kts @@ -2,6 +2,7 @@ import org.jetbrains.kotlin.gradle.dsl.JvmTarget import org.jetbrains.kotlin.gradle.targets.native.tasks.KotlinNativeSimulatorTest +import java.util.Properties plugins { alias(libs.plugins.kotlinMultiplatform) @@ -10,6 +11,13 @@ plugins { alias(libs.plugins.mavenPublish) } +val conformancePins = Properties().apply { + rootProject.file("test-fixtures/torrent/clients.properties").inputStream().use { load(it) } +} +val libtorrentFixtureVersion = conformancePins.getProperty("libtorrent4j.version") +val torrentConformance = providers.gradleProperty("torrentConformance") + .map(String::toBooleanStrict).orElse(false) + kotlin { android { withHostTest {} @@ -58,11 +66,12 @@ kotlin { implementation(projects.library.server) implementation(projects.library.remote) implementation(libs.ktor.client.cio) - implementation("org.libtorrent4j:libtorrent4j:2.1.0-39") - runtimeOnly("org.libtorrent4j:libtorrent4j-macos:2.1.0-39") - runtimeOnly("org.libtorrent4j:libtorrent4j-linux:2.1.0-39") - runtimeOnly("org.libtorrent4j:libtorrent4j-windows:2.1.0-39") + implementation("org.libtorrent4j:libtorrent4j:$libtorrentFixtureVersion") + runtimeOnly("org.libtorrent4j:libtorrent4j-macos:$libtorrentFixtureVersion") + runtimeOnly("org.libtorrent4j:libtorrent4j-linux:$libtorrentFixtureVersion") + runtimeOnly("org.libtorrent4j:libtorrent4j-windows:$libtorrentFixtureVersion") } + jvmTest.get().resources.srcDir(rootProject.file("test-fixtures/torrent")) named("androidDeviceTest") { dependencies { implementation(libs.kotlin.test) @@ -83,8 +92,21 @@ tasks.withType().configureEach { // Explicit opt-in inputs make external-client and package smoke runs reproducible under Gradle. tasks.withType().configureEach { + val conformance = torrentConformance.get() + inputs.property("torrentConformance", conformance) + // Required evidence must describe this execution, not restored/cached XML from an earlier run. + val benchmark = providers.environmentVariable("KETCH_TORRENT_BENCHMARK").orNull == "1" + outputs.upToDateWhen { !conformance && !benchmark } + outputs.cacheIf { !conformance && !benchmark } + if (!conformance && providers.environmentVariable("TRANSMISSION_DAEMON").orNull.isNullOrBlank()) { + filter.excludeTestsMatching("*TransmissionInteropTest") + } + if (providers.environmentVariable("KETCH_TORRENT_BENCHMARK").orNull != "1") { + filter.excludeTestsMatching("*TorrentBenchmarkTest") + } for (name in listOf("TRANSMISSION_DAEMON", "KETCH_TORRENT_BENCHMARK", - "KETCH_NATIVE_CLI", "KETCH_JVM_CLI")) { + "KETCH_NATIVE_CLI", "KETCH_JVM_CLI", "KETCH_BENCHMARK_BYTES", "KETCH_BENCHMARK_RUNS", + "KETCH_BENCHMARK_REVISION", "KETCH_BENCHMARK_REPORT")) { val value = providers.environmentVariable(name).orElse("") inputs.property(name, value) environment(name, value.get()) 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..0a44d7ce7 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 @@ -4,6 +4,7 @@ import kotlinx.coroutines.CancellationException import kotlinx.coroutines.CoroutineStart import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.Job import kotlinx.coroutines.NonCancellable import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.async @@ -13,10 +14,12 @@ 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 import kotlinx.coroutines.supervisorScope +import kotlinx.coroutines.sync.Semaphore import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.withLock import kotlinx.coroutines.withContext @@ -40,14 +43,20 @@ 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 cache = TorrentMetadataCache(scope) + private val exchangeBudgets = TorrentExchangeBudgets(config) + private val budget = exchangeBudgets.transfer + private val storageSlots = Semaphore(config.maxOpenPayloadFiles) + private val admissions = TorrentAdmissionLedger(exchangeBudgets.sessions) + internal val admittedSessionBytes: Int get() = exchangeBudgets.sessions.allocated + private val cache = TorrentMetadataCache(scope, config.maxCachedMetadataBytes, + exchangeBudgets.cache) private val tracker = TorrentTracker(http, network) private val peerId = torrentRandomBytes(20) private val mutex = Mutex() private val dhtMutex = Mutex() private val sessions = mutableMapOf() private val outputs = mutableMapOf() + private val sessionLeases = mutableMapOf() private val running = AtomicBoolean(false) private var closed = false private var port = 0 @@ -57,6 +66,10 @@ internal class KotlinTorrentEngine( private val uploadRate = TorrentRateLimiter() override val isRunning: Boolean get() = running.load() + init { + checkNotNull(scope.coroutineContext[Job]).invokeOnCompletion { admissions.close() } + } + override suspend fun start() = mutex.withLock { check(!closed) { "Torrent runtime is closed" } if (running.load()) return@withLock @@ -101,6 +114,10 @@ internal class KotlinTorrentEngine( } } } + } catch (error: Exception) { + // Closing a socket during shutdown can throw before its provider observes cancellation. + currentCoroutineContext().ensureActive() + throw error } finally { listener.close() } } @@ -116,7 +133,7 @@ internal class KotlinTorrentEngine( if (closed) return closed = true running.store(false) - sessions.values.toList().also { sessions.clear(); outputs.clear() } + sessions.values.toList().also { sessions.clear(); outputs.clear(); sessionLeases.clear() } } withContext(NonCancellable) { try { @@ -124,8 +141,12 @@ internal class KotlinTorrentEngine( nodes?.forEach { it.close() } } finally { scope.cancel() - network.close() - http.close() + try { network.close() } finally { + try { http.close() } finally { + cache.close() + checkNotNull(scope.coroutineContext[Job]).join() + } + } } } } @@ -146,7 +167,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) { @@ -212,28 +233,39 @@ internal class KotlinTorrentEngine( check(sessions.size < config.maxActiveTorrents) { "Too many active torrents" } val hash = spec.metadata.infoHash.hex check(hash !in sessions) { "Torrent already has an active owner" } - val requested = FileSystem.SYSTEM.canonicalize(".".toPath()) - .resolve(spec.outputPath).normalized() - val store = TorrentPieceStore(spec.metadata, requested, spec.selected, spec.taskId) - 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" } - val checkpoint = spec.resumeData?.let(TorrentCheckpoint::decode) - val session = KotlinTorrentSession(store, network, budget, scope, - connections = config.connectionsPerTorrent, uploadPolicy = config.effectiveUploadPolicy, - checkpoint = checkpoint, peerId = peerId, - discover = { peers, owner -> discover(spec, peers, owner) }, - downloadThrottle = { downloadRate.acquire(it); spec.throttle(it) }, - uploadThrottle = { uploadRate.acquire(it) }, - ) - sessions[hash] = session - outputs[hash] = output.toString() - session + 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" } + val checkpoint = spec.resumeData?.let(TorrentCheckpoint::decode) + val session = KotlinTorrentSession(store, network, budget, scope, + connections = config.connectionsPerTorrent, uploadPolicy = config.effectiveUploadPolicy, + checkpoint = checkpoint, peerId = peerId, + discover = { peers, owner -> discover(spec, peers, owner) }, + downloadThrottle = { downloadRate.acquire(it); spec.throttle(it) }, + uploadThrottle = { uploadRate.acquire(it) }, + ) + sessionLeases[hash] = lease + sessions[hash] = session + outputs[hash] = output.toString() + session + } catch (failure: Throwable) { + sessions.remove(hash) + outputs.remove(hash) + sessionLeases.remove(hash) + admissions.release(lease) + throw failure + } } private suspend fun discover( @@ -334,6 +366,7 @@ internal class KotlinTorrentEngine( sessions[infoHash]?.close(deleteFiles) sessions.remove(infoHash) outputs.remove(infoHash) + sessionLeases.remove(infoHash)?.let(admissions::release) } } diff --git a/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/KotlinTorrentSession.kt b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/KotlinTorrentSession.kt index fd59746e8..4bec74b45 100644 --- a/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/KotlinTorrentSession.kt +++ b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/KotlinTorrentSession.kt @@ -79,6 +79,7 @@ internal class KotlinTorrentSession( private val lastPayload = AtomicLong(0) private var job: Job? = null private var closed = false + private var filesDeleted = false private var recovered = false private val _state = MutableStateFlow(TorrentSessionState.PAUSED) private val _downloadedBytes = MutableStateFlow(0L) @@ -203,18 +204,20 @@ internal class KotlinTorrentSession( } suspend fun close(deleteFiles: Boolean = false) = lifecycle.withLock { - if (closed) return@withLock - closed = true - job?.cancelAndJoin() - job = null - scope.cancel() - incoming.cancel() - resets.cancel() - _state.value = TorrentSessionState.STOPPED - if (deleteFiles) { + if (!closed) { + job?.cancelAndJoin() + job = null + scope.cancel() + incoming.cancel() + resets.cancel() + closed = true + _state.value = TorrentSessionState.STOPPED + } + if (deleteFiles && !filesDeleted) { if (!recovered) checkpoint?.let { store.restore(it) } store.recoverOwnership() store.cleanup() + filesDeleted = true } } } 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..16738f46c --- /dev/null +++ b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentBufferBudget.kt @@ -0,0 +1,52 @@ +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 cache = TorrentBufferBudget(config.maxCachedMetadataBytes, root) + val sessions = TorrentBufferBudget(config.maxSessionStateBytes, 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..4132df807 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,21 @@ data class TorrentConfig( val maxMetadataBytes: Int = 4 * 1024 * 1024, /** Maximum simultaneously buffered piece data across the engine. */ val maxBufferedBytes: Int = 32 * 1024 * 1024, + /** + * Combined admission ceiling for buffers, metadata exchange, cache entries, and session state. + * This is not a process RSS limit; platform allocations and caller-owned state are separate. + */ + val maxExchangeBytes: Int = 64 * 1024 * 1024, + /** Cache retention ceiling, including conservative file/index and string allowances. */ + val maxCachedMetadataBytes: Int = 4 * 1024 * 1024, + /** Aggregate admission allowance for session metadata, indexes, and checking scratch space. */ + val maxSessionStateBytes: Int = 8 * 1024 * 1024, + /** Aggregate open payload-file ceiling. Storage waits before opening another payload handle. */ + val maxOpenPayloadFiles: Int = 32, + /** File count ceiling applied before constructing a session's storage indexes. */ + val maxFilesPerTorrent: Int = 10_000, + /** Logical piece ceiling applied before constructing session and scheduler arrays. */ + val maxPiecesPerTorrent: Int = 250_000, ) { init { @@ -47,8 +62,21 @@ 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(maxCachedMetadataBytes > 0 && maxSessionStateBytes > 0) + require(maxOpenPayloadFiles in 1..128) + require(maxFilesPerTorrent in 1..100_000) + require(maxPiecesPerTorrent in 1..1_000_000) + require(maxBufferedBytes.toLong() + metadataExchangeBytes + maxCachedMetadataBytes + + maxSessionStateBytes <= + maxExchangeBytes.toLong()) { + "maxExchangeBytes must cover transfers, metadata exchange, cache, and session state" + } } + /** 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/TorrentMetadataCache.kt b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentMetadataCache.kt index 2f4c19e1c..5c4731308 100644 --- a/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentMetadataCache.kt +++ b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentMetadataCache.kt @@ -1,10 +1,15 @@ package com.linroid.ketch.torrent import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.CoroutineStart import kotlinx.coroutines.Deferred +import kotlinx.coroutines.Job import kotlinx.coroutines.NonCancellable import kotlinx.coroutines.async +import kotlinx.coroutines.awaitCancellation +import kotlinx.coroutines.cancelAndJoin import kotlinx.coroutines.ensureActive +import kotlinx.coroutines.launch import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.withLock import kotlinx.coroutines.withContext @@ -12,24 +17,37 @@ import kotlinx.coroutines.withContext /** Runtime-owned shared fetches; cancellation of one caller does not cancel other waiters. */ internal class TorrentMetadataCache( private val scope: CoroutineScope, - private val capacityBytes: Int = 32 * 1024 * 1024, + capacityBytes: Int = 32 * 1024 * 1024, + parentBudget: TorrentBufferBudget? = null, ) { - init { require(capacityBytes > 0) } + private val budget = TorrentBufferBudget(capacityBytes, parentBudget) + internal val retainedBytes: Int get() = budget.allocated private val mutex = Mutex() - private val entries = linkedMapOf() + private class Entry(val metadata: TorrentMetadata, val lease: TorrentBufferBudget.Lease) + private val entries = linkedMapOf() private val pending = mutableMapOf>() - private var bytes = 0L + private var closed = false + private var cleanupJob: Job? = null + + init { + // Start before returning the cache so even immediate owner cancellation runs cleanup. + cleanupJob = scope.launch(start = CoroutineStart.UNDISPATCHED) { + try { awaitCancellation() } finally { withContext(NonCancellable) { close() } } + } + } suspend fun get(hash: InfoHash): TorrentMetadata? = mutex.withLock { scope.coroutineContext.ensureActive() - entries.remove(hash)?.also { entries[hash] = it } + check(!closed) { "Metadata cache is closed" } + entries.remove(hash)?.also { entries[hash] = it }?.metadata } suspend fun resolve(hash: InfoHash, fetch: suspend () -> TorrentMetadata): TorrentMetadata { get(hash)?.let { return it } val operation = mutex.withLock { scope.coroutineContext.ensureActive() + check(!closed) { "Metadata cache is closed" } pending.getOrPut(hash) { check(pending.size < 16) { "Too many pending metadata requests" } scope.async { @@ -47,27 +65,56 @@ internal class TorrentMetadataCache( } suspend fun put(value: TorrentMetadata): TorrentMetadata = mutex.withLock { - // The info hash authenticates only the info dictionary, never caller tracker credentials. - val metadata = value.copy( - trackers = emptyList(), - trackerTiers = emptyList(), - comment = null, - createdBy = null, - metainfoBytes = metainfoFromInfo(value.infoBytes, emptyList()), - ) scope.coroutineContext.ensureActive() - val size = weight(metadata) - entries.remove(metadata.infoHash)?.let { bytes -= weight(it) } - if (size > capacityBytes) return@withLock metadata - while (entries.isNotEmpty() && (bytes + size > capacityBytes || entries.size >= 8)) { - val key = entries.keys.first() - bytes -= weight(checkNotNull(entries.remove(key))) + check(!closed) { "Metadata cache is closed" } + val size = cacheWeight(value) + entries.remove(value.infoHash)?.lease?.close() + var lease: TorrentBufferBudget.Lease? = null + if (size <= budget.capacity) { + if (entries.size >= 8) evictOldest() + lease = budget.reserve(size.toInt()) + while (lease == null && entries.isNotEmpty()) { + evictOldest() + lease = budget.reserve(size.toInt()) + } + } + try { + // The hash authenticates the info dictionary, never caller tracker credentials. + val metadata = value.copy( + trackers = emptyList(), + trackerTiers = emptyList(), + comment = null, + createdBy = null, + metainfoBytes = metainfoFromInfo(value.infoBytes, emptyList()), + ) + if (lease != null) entries[metadata.infoHash] = Entry(metadata, lease) + metadata + } catch (failure: Throwable) { + lease?.close() + throw failure } - entries[metadata.infoHash] = metadata - bytes += size - metadata } - private fun weight(metadata: TorrentMetadata): Long = metadata.metainfoBytes.size.toLong() + - metadata.infoBytes.size + metadata.pieceHashes.size + metadata.files.sumOf { it.path.length * 2L } + /** Reject future use, release retained entries, and await canceled shared fetches. */ + suspend fun close(): Unit = withContext(NonCancellable) { + cleanupJob?.cancel() + val operations = mutex.withLock { + closed = true + entries.values.forEach { it.lease.close() } + entries.clear() + pending.values.toList() + } + // Fetch finalizers take the same mutex, so never await them while holding it. + operations.forEach { it.cancel() } + operations.forEach { it.cancelAndJoin() } + } + + private fun evictOldest() { + entries.remove(entries.keys.first())?.lease?.close() + } } + +/** Retained arrays, strings, file records, and container allowance; not process RSS. */ +internal fun cacheWeight(metadata: TorrentMetadata): Long = + metadata.infoBytes.size.toLong() * 2 + metadata.pieceHashes.size + + metadata.name.length * 2L + metadata.files.sumOf { it.path.length * 2L + 64 } + 512 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/commonMain/kotlin/com/linroid/ketch/torrent/TorrentPieceStore.kt b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentPieceStore.kt index 1e0679a8d..080006680 100644 --- a/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentPieceStore.kt +++ b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentPieceStore.kt @@ -1,10 +1,13 @@ package com.linroid.ketch.torrent +import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.IO import kotlinx.coroutines.currentCoroutineContext import kotlinx.coroutines.ensureActive import okio.EOFException +import kotlinx.coroutines.sync.Semaphore +import kotlinx.coroutines.sync.withPermit import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.withLock import kotlinx.coroutines.withContext @@ -20,6 +23,7 @@ internal class TorrentPieceStore( selected: Set, private val taskId: String, private val fileSystem: FileSystem = torrentFileSystem, + private val storageSlots: Semaphore = Semaphore(32), ) { private val mutex = Mutex() private val output = absolute(output) @@ -90,8 +94,8 @@ internal class TorrentPieceStore( fun needed(index: Int): Boolean = wanted[index] suspend fun initialize() = mutex.withLock { - withContext(Dispatchers.IO) { - if (initialized) return@withContext + storageOperation { + if (initialized) return@storageOperation ensureDirectory(sidecar) loadOwnership() journal.initialize() @@ -128,7 +132,7 @@ internal class TorrentPieceStore( require(bytes.size == pieceSize(index) && needed(index)) if (!matches(index, bytes)) return@withLock false if (verified[index]) return@withLock true - withContext(Dispatchers.IO) { + storageOperation { val start = index * metadata.pieceLength val end = start + bytes.size val spans = overlappingFiles(index).filter { it in selected } @@ -155,12 +159,12 @@ internal class TorrentPieceStore( /** Reads selected spans from the output and skipped boundary spans from owned sidecars. */ suspend fun read(index: Int): ByteArray = mutex.withLock { - withContext(Dispatchers.IO) { readPiece(index) } + storageOperation { readPiece(index) } } suspend fun recheck(): BooleanArray = mutex.withLock { check(initialized) - withContext(Dispatchers.IO) { + storageOperation { verified.fill(false) fileProgress.fill(0) remaining = wanted.count { it } @@ -194,7 +198,7 @@ internal class TorrentPieceStore( suspend fun finish() = mutex.withLock { check(initialized && remaining == 0) { "Selected torrent files are incomplete" } - withContext(Dispatchers.IO) { + storageOperation { for (file in selected) { val path = filePath(file) validateRegularPath(path) @@ -212,7 +216,7 @@ internal class TorrentPieceStore( /** Deletes only recorded paths whose OS identity still matches the task-owned file. */ suspend fun cleanup() = mutex.withLock { - withContext(Dispatchers.IO) { + storageOperation { val files = ownedFiles.toList().asReversed().sortedBy { it.first == journalPath } for ((path, identity) in files) { validateRegularPath(path) @@ -235,7 +239,7 @@ internal class TorrentPieceStore( } suspend fun recoverOwnership() = mutex.withLock { - withContext(Dispatchers.IO) { loadOwnership() } + storageOperation { loadOwnership() } } private fun loadOwnership() { @@ -261,7 +265,7 @@ internal class TorrentPieceStore( mutex.withLock { require(receivedBytes >= 0 && uploadedBytes >= 0) check(initialized) - withContext(Dispatchers.IO) { + storageOperation { if (journal.needsCompaction) { require(torrentFileIdentity(journalPath) == ownedFiles[journalPath]) val records = ownedDirectories.map { true to TorrentOwnedPath(it.key.toString(), it.value) } + @@ -291,6 +295,12 @@ internal class TorrentPieceStore( } } + // A store opens at most one payload handle at a time. Acquire before dispatching blocking I/O; + // retain the slot until that I/O actually returns, including cancellation and provider failures. + // Journals/checkpoints may open extra non-payload handles and are outside this payload ceiling. + private suspend fun storageOperation(block: suspend CoroutineScope.() -> T): T = + storageSlots.withPermit { withContext(Dispatchers.IO, block) } + private fun snapshot(): TorrentCheckpoint = TorrentCheckpoint(taskId, metadata, output.toString(), selected, verified.copyOf(), ownedFiles.map { TorrentOwnedPath(it.key.toString(), it.value) }, diff --git a/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentSessionAdmission.kt b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentSessionAdmission.kt new file mode 100644 index 000000000..2ba85eb9a --- /dev/null +++ b/library/torrent/src/commonMain/kotlin/com/linroid/ketch/torrent/TorrentSessionAdmission.kt @@ -0,0 +1,70 @@ +package com.linroid.ketch.torrent + +import kotlin.concurrent.atomics.AtomicReference +import kotlin.concurrent.atomics.ExperimentalAtomicApi + +/** Apply every limit together before storage construction or checkpoint decoding allocates indexes. */ +internal fun admitSession( + spec: TorrentTaskSpec, + config: TorrentConfig, + budget: TorrentBufferBudget, +): TorrentBufferBudget.Lease { + require(spec.metadata.files.size <= config.maxFilesPerTorrent) { "Torrent file limit exceeded" } + require(spec.metadata.pieceHashes.size / 20 <= config.maxPiecesPerTorrent) { + "Torrent piece limit exceeded" + } + val bytes = sessionStateWeight(spec) + check(bytes <= budget.capacity) { "Torrent session state exceeds admission capacity" } + return checkNotNull(budget.reserve(bytes.toInt())) { "Torrent session state budget exhausted" } +} + +/** + * Conservative allowance for retained metadata, storage/scheduler arrays and boxed candidate lists, + * path/ownership indexes, two checking buffers, and checkpoint decoding. Peer bitfields and wire + * buffers are already charged by TorrentSwarm. This is admission accounting, not measured heap/RSS. + * Keep arithmetic in Long; an over-capacity estimate must fail before narrowing to an Int. + */ +internal fun sessionStateWeight(spec: TorrentTaskSpec): Long { + val metadata = spec.metadata + val pieces = metadata.pieceHashes.size.toLong() / 20 + val largestPiece = minOf(metadata.pieceLength, metadata.totalBytes) + require(largestPiece >= 0) + return cacheWeight(metadata) + metadata.metainfoBytes.size + + metadata.trackers.sumOf { it.length * 4L } + + (metadata.comment?.length ?: 0) * 4L + (metadata.createdBy?.length ?: 0) * 4L + + pieces * 128 + metadata.files.sumOf { 512 + it.path.length * 4L } + + largestPiece * 2 + (spec.resumeData?.size ?: 0) * 8L + + spec.outputPath.length * 4L + (spec.magnetUri?.length ?: 0) * 4L + + spec.selected.size * 64L + 128 * 1024 +} + +/** Runtime ownership ledger; a stopped-but-registered session remains charged after cleanup failure. */ +@OptIn(ExperimentalAtomicApi::class) +internal class TorrentAdmissionLedger(private val budget: TorrentBufferBudget) { + private val entries = AtomicReference?>(emptyList()) + + fun admit(spec: TorrentTaskSpec, config: TorrentConfig): TorrentBufferBudget.Lease { + val lease = admitSession(spec, config, budget) + while (true) { + val current = entries.load() + if (current == null) { + lease.close() + error("Torrent admission is closed") + } + if (entries.compareAndSet(current, current + lease)) return lease + } + } + + fun release(lease: TorrentBufferBudget.Lease) { + while (true) { + val current = entries.load() ?: break + if (lease !in current || entries.compareAndSet(current, current.filter { it !== lease })) break + } + lease.close() + } + + /** Called only after the owning runtime's jobs have finished. No new admission is possible. */ + fun close() { + entries.exchange(null)?.forEach { it.close() } + } +} 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..b81556158 --- /dev/null +++ b/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentBufferBudgetTest.kt @@ -0,0 +1,74 @@ +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 cached = assertNotNull(budgets.cache.reserve(budgets.cache.capacity)) + val state = assertNotNull(budgets.sessions.reserve(budgets.sessions.capacity)) + val info = assertNotNull(metadata.reserve(metadata.capacity)) + assertNull(transfer.reserve(1)) + assertEquals(transfer.capacity + metadata.capacity + budgets.cache.capacity + + budgets.sessions.capacity, budgets.allocated) + info.close() + payload.close() + cached.close() + state.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/TorrentListenerShutdownTest.kt b/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentListenerShutdownTest.kt new file mode 100644 index 000000000..10c3cab94 --- /dev/null +++ b/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentListenerShutdownTest.kt @@ -0,0 +1,49 @@ +package com.linroid.ketch.torrent + +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.withContext +import kotlinx.coroutines.CompletableDeferred +import kotlinx.coroutines.awaitCancellation +import kotlinx.coroutines.test.runTest +import kotlinx.coroutines.withTimeout +import okio.IOException +import kotlin.test.Test +import kotlin.test.assertTrue + +class TorrentListenerShutdownTest { + @Test + fun socketFailureDuringCanceledAcceptDoesNotEscapeEngineShutdown() = runTest { + val entered = CompletableDeferred() + val closed = CompletableDeferred() + val network = object : TorrentNetwork { + override suspend fun connect(remote: PeerEndpoint): TorrentConnection = error("Unused") + override suspend fun bindUdp(local: PeerEndpoint): TorrentDatagramSocket = error("Unused") + override suspend fun listen(local: PeerEndpoint): TorrentListener { + if (local.host == "::") throw IOException("No second listener") + return object : TorrentListener { + override val local = PeerEndpoint("127.0.0.1", 12345) + override suspend fun accept(): TorrentConnection { + entered.complete(Unit) + try { awaitCancellation() } finally { + // Socket providers can report closure instead of coroutine cancellation. + throw IOException("Listener closed during accept") + } + } + override fun close() { closed.complete(Unit) } + } + } + override fun close() = Unit + } + val engine = KotlinTorrentEngine(TorrentConfig(dhtEnabled = false), rawNetwork = network) + try { + withContext(Dispatchers.Default) { + withTimeout(5000) { + engine.start() + entered.await() + engine.stop() + assertTrue(closed.isCompleted) + } + } + } finally { engine.stop() } + } +} diff --git a/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentMetadataCacheTest.kt b/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentMetadataCacheTest.kt index 2e1ec7d05..45a3f0650 100644 --- a/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentMetadataCacheTest.kt +++ b/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentMetadataCacheTest.kt @@ -1,12 +1,22 @@ package com.linroid.ketch.torrent import kotlinx.coroutines.CompletableDeferred +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Job +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.awaitCancellation import kotlinx.coroutines.async import kotlinx.coroutines.cancelAndJoin +import kotlinx.coroutines.coroutineScope +import kotlinx.coroutines.withTimeout import kotlinx.coroutines.test.runCurrent import kotlinx.coroutines.test.runTest import kotlin.test.Test import kotlin.test.assertEquals +import kotlin.test.assertFailsWith +import kotlin.test.assertNotNull +import kotlin.test.assertNull +import kotlin.test.assertTrue class TorrentMetadataCacheTest { @Test @@ -54,4 +64,101 @@ class TorrentMetadataCacheTest { assertEquals(metadata.infoHash, cache.resolve(metadata.infoHash, fetch).infoHash) assertEquals(1, fetches) } + + @Test + fun evictionReplacementAndCloseReleaseParentReservations() = runTest { + val first = fixture("a") + val second = fixture("b") + val third = fixture("c") + val size = cacheWeight(first).toInt() + val root = TorrentBufferBudget(size * 2) + val cache = TorrentMetadataCache(backgroundScope, size * 2, root) + cache.put(first) + cache.put(second) + assertNotNull(cache.get(first.infoHash)) // Refresh first; second is now least recently used. + cache.put(third) + assertNull(cache.get(second.infoHash)) + assertNotNull(cache.get(first.infoHash)) + cache.put(third) // Replacement must not leak the earlier reservation. + assertEquals(size * 2, cache.retainedBytes) + assertEquals(size * 2, root.allocated) + cache.close() + cache.close() + assertEquals(0, root.allocated) + assertFailsWith { cache.put(first) } + assertFailsWith { cache.get(first.infoHash) } + } + + @Test + fun overBudgetResultsAreReturnedWithoutRetainingAnEntry() = runTest { + val metadata = fixture("a") + val size = cacheWeight(metadata).toInt() + val root = TorrentBufferBudget(size) + val cache = TorrentMetadataCache(backgroundScope, size, root) + val other = assertNotNull(root.reserve(size)) + assertEquals(metadata.infoHash, cache.put(metadata).infoHash) + assertNull(cache.get(metadata.infoHash)) + assertEquals(0, cache.retainedBytes) + other.close() + cache.put(metadata) + assertEquals(size, root.allocated) + cache.close() + val small = TorrentMetadataCache(backgroundScope, size - 1, root) + assertEquals(metadata.infoHash, small.put(metadata).infoHash) + assertNull(small.get(metadata.infoHash)) + assertEquals(0, root.allocated) + } + + @Test + fun closeCancelsSharedFetchAndWaitsForItsFinalizer() = runTest { + val metadata = fixture("a") + val cache = TorrentMetadataCache(backgroundScope) + val ready = CompletableDeferred() + var finalized = false + val caller = async { + cache.resolve(metadata.infoHash) { + ready.complete(Unit) + try { awaitCancellation() } finally { finalized = true } + } + } + ready.await() + cache.close() + assertTrue(finalized) + caller.join() + assertTrue(caller.isCancelled) + assertEquals(0, cache.retainedBytes) + } + + @Test + fun ownerCancellationReleasesRetainedMetadata() = runTest { + val ownerJob = SupervisorJob() + val owner = CoroutineScope(backgroundScope.coroutineContext + ownerJob) + val root = TorrentBufferBudget(4096) + val cache = TorrentMetadataCache(owner, 4096, root) + try { + cache.put(fixture("a")) + assertTrue(root.allocated > 0) + } finally { + owner.coroutineContext[Job]!!.cancelAndJoin() + } + assertEquals(0, root.allocated) + } + + @Test + fun explicitCloseAllowsStructuredOwnerToComplete() = runTest { + withTimeout(1000) { + coroutineScope { + val cache = TorrentMetadataCache(this) + cache.put(fixture("a")) + cache.close() + } + } + } + + private fun fixture(name: String): TorrentMetadata = TorrentMetadata.fromBencode( + Bencode.encode(mapOf("info" to mapOf( + "name" to name, "length" to 0L, "piece length" to 16_384L, "pieces" to ByteArray(0) + ))) + ) + } 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) } } } diff --git a/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentPieceStoreTest.kt b/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentPieceStoreTest.kt index 41826f263..166c9a179 100644 --- a/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentPieceStoreTest.kt +++ b/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentPieceStoreTest.kt @@ -1,6 +1,10 @@ package com.linroid.ketch.torrent import kotlinx.coroutines.CancellationException +import kotlinx.coroutines.CoroutineStart +import kotlinx.coroutines.async +import kotlinx.coroutines.cancelAndJoin +import kotlinx.coroutines.sync.Semaphore import kotlinx.coroutines.test.runTest import okio.FileSystem import okio.FileHandle @@ -46,6 +50,54 @@ class TorrentPieceStoreTest { } } + @Test + fun canceledStorageWaitDoesNotTouchDiskOrConsumeSharedSlot() = runTest { + val slots = Semaphore(1) + val root = FileSystem.SYSTEM_TEMPORARY_DIRECTORY / + "ketch-storage-wait-${InfoHash.fromBytes(torrentRandomBytes(20)).hex}" + val waiting = TorrentPieceStore(metadata, root / "first", emptySet(), "first", + storageSlots = slots) + val next = TorrentPieceStore(metadata, root / "second", emptySet(), "second", + storageSlots = slots) + try { + slots.acquire() + val initialization = async(start = CoroutineStart.UNDISPATCHED) { waiting.initialize() } + assertFalse(initialization.isCompleted) + assertFalse(torrentFileSystem.exists(root)) + initialization.cancelAndJoin() + slots.release() + next.initialize() + assertTrue(next.isInitialized()) + assertFalse(torrentFileSystem.exists(root / "first")) + assertEquals(1, slots.availablePermits) + } finally { + torrentFileSystem.deleteRecursively(root, mustExist = false) + } + } + + @Test + fun failedStorageOperationReturnsSlotToOtherTorrent() = runTest { + val slots = Semaphore(1) + val root = FileSystem.SYSTEM_TEMPORARY_DIRECTORY / + "ketch-storage-failure-${InfoHash.fromBytes(torrentRandomBytes(20)).hex}" + val failing = object : ForwardingFileSystem(torrentFileSystem) { + override fun openReadWrite(file: Path, mustCreate: Boolean, mustExist: Boolean): FileHandle { + throw IOException("Provider unavailable") + } + } + val first = TorrentPieceStore(metadata, root / "first", emptySet(), "first", failing, slots) + val next = TorrentPieceStore(metadata, root / "second", emptySet(), "second", + storageSlots = slots) + try { + assertFailsWith { first.initialize() } + assertEquals(1, slots.availablePermits) + next.initialize() + assertTrue(next.isInitialized()) + } finally { + torrentFileSystem.deleteRecursively(root, mustExist = false) + } + } + private val bytes = "0123456789".encodeToByteArray() private val metadata = TorrentMetadata.fromBencode(Bencode.encode(mapOf("info" to mapOf( "name" to "pack", "piece length" to 4L, diff --git a/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentSessionAdmissionTest.kt b/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentSessionAdmissionTest.kt new file mode 100644 index 000000000..c6bd5b249 --- /dev/null +++ b/library/torrent/src/commonTest/kotlin/com/linroid/ketch/torrent/TorrentSessionAdmissionTest.kt @@ -0,0 +1,148 @@ +package com.linroid.ketch.torrent + +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.delay +import kotlinx.coroutines.test.runTest +import kotlinx.coroutines.withContext +import kotlinx.coroutines.withTimeout +import okio.FileSystem +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertFailsWith +import kotlin.test.assertFalse +import kotlin.test.assertTrue + +class TorrentSessionAdmissionTest { + @Test + fun countAndMemoryLimitsApplyTogetherWithoutConsumingCreditOnFailure() { + val metadata = fixture("a") + val spec = TorrentTaskSpec("task", metadata, "/tmp/admission/a", emptySet()) + val size = sessionStateWeight(spec).toInt() + val budget = TorrentBufferBudget(size) + val config = TorrentConfig(maxPiecesPerTorrent = 1, maxFilesPerTorrent = 1) + assertFailsWith { + admitSession(spec.copy(metadata = metadata.copy(pieceHashes = ByteArray(40))), config, budget) + } + assertFailsWith { + admitSession(spec.copy(metadata = metadata.copy(files = metadata.files + metadata.files)), + config, budget) + } + assertFailsWith { + admitSession(spec.copy(resumeData = ByteArray(size)), config, budget) + } + assertEquals(0, budget.allocated) + val lease = admitSession(spec, config, budget) + assertEquals(size, budget.allocated) + assertFailsWith { admitSession(spec, config, budget) } + lease.close() + assertEquals(0, budget.allocated) + } + + @Test + fun engineRejectsBeforeCreatingStorageAndReturnsCreditOnFailedConstruction() = runTest { + withContext(Dispatchers.Default) { + val root = FileSystem.SYSTEM_TEMPORARY_DIRECTORY / + "admission-${InfoHash.fromBytes(torrentRandomBytes(20)).hex}" + val metadata = fixture("a") + val spec = TorrentTaskSpec("valid", metadata, (root / "payload").toString(), emptySet()) + val size = sessionStateWeight(spec).toInt() + val engine = KotlinTorrentEngine(TorrentConfig(dhtEnabled = false, maxSessionStateBytes = size)) + try { + engine.start() + assertFailsWith { engine.addTask(spec.copy(taskId = "!")) } + assertEquals(0, engine.admittedSessionBytes) + assertFalse(FileSystem.SYSTEM.exists(root)) + val session = engine.addTask(spec) + assertEquals(size, engine.admittedSessionBytes) + val second = spec.copy(taskId = "other", metadata = fixture("b")) + assertFailsWith { engine.addTask(second) } + assertFalse(FileSystem.SYSTEM.exists(root)) + session.pause() + assertEquals(size, engine.admittedSessionBytes) // Paused sessions still own their indexes. + engine.removeTorrent(metadata.infoHash.hex, false) + assertEquals(0, engine.admittedSessionBytes) + engine.addTask(second) + assertTrue(engine.admittedSessionBytes > 0) + } finally { + engine.stop() + FileSystem.SYSTEM.deleteRecursively(root, mustExist = false) + } + assertEquals(0, engine.admittedSessionBytes) + } + } + + @Test + fun nonSuspendingEngineCloseEventuallyReturnsSessionCredit() = runTest { + withContext(Dispatchers.Default) { + val engine = KotlinTorrentEngine(TorrentConfig(dhtEnabled = false)) + try { + engine.start() + engine.addTask(TorrentTaskSpec("task", fixture("a"), + (FileSystem.SYSTEM_TEMPORARY_DIRECTORY / "admission-close").toString(), emptySet())) + assertTrue(engine.admittedSessionBytes > 0) + engine.close() + withTimeout(5000) { + while (engine.admittedSessionBytes != 0) delay(1) + } + } finally { engine.stop() } + } + } + + @Test + fun failedDeletionRemainsChargedAndCannotSilentlySucceedOnRetry() = runTest { + withContext(Dispatchers.Default) { + val root = FileSystem.SYSTEM_TEMPORARY_DIRECTORY / + "admission-removal-${InfoHash.fromBytes(torrentRandomBytes(20)).hex}" + val metadata = fixture("a") + val output = (root / "payload").toString() + val checkpoint = TorrentCheckpoint("different-task", metadata, output, emptySet(), + BooleanArray(1), emptyList(), emptyList()).encode() + val spec = TorrentTaskSpec("task", metadata, output, emptySet(), resumeData = checkpoint) + val size = sessionStateWeight(spec).toInt() + val engine = KotlinTorrentEngine(TorrentConfig(dhtEnabled = false, maxSessionStateBytes = size)) + try { + engine.start() + engine.addTask(spec) + repeat(2) { + assertFailsWith { + engine.removeTorrent(metadata.infoHash.hex, deleteFiles = true) + } + assertEquals(size, engine.admittedSessionBytes) + } + val other = TorrentTaskSpec("other", fixture("b"), (root / "other").toString(), emptySet()) + assertFailsWith { engine.addTask(other) } + assertFalse(FileSystem.SYSTEM.exists(root)) + engine.removeTorrent(metadata.infoHash.hex, deleteFiles = false) + assertEquals(0, engine.admittedSessionBytes) + engine.addTask(other) + } finally { + engine.stop() + FileSystem.SYSTEM.deleteRecursively(root, mustExist = false) + } + assertEquals(0, engine.admittedSessionBytes) + } + } + + @Test + fun closedRuntimeLedgerRejectsNewAdmissionAndReturnsCreditOnce() { + val config = TorrentConfig() + val budget = TorrentBufferBudget(config.maxSessionStateBytes) + val ledger = TorrentAdmissionLedger(budget) + val spec = TorrentTaskSpec("task", fixture("a"), "/tmp/admission-ledger", emptySet()) + val lease = ledger.admit(spec, config) + assertTrue(budget.allocated > 0) + ledger.close() + ledger.release(lease) + ledger.close() + assertEquals(0, budget.allocated) + assertFailsWith { ledger.admit(spec, config) } + assertEquals(0, budget.allocated) + } + + private fun fixture(name: String): TorrentMetadata = TorrentMetadata.fromBencode( + Bencode.encode(mapOf("info" to mapOf( + "name" to name, "length" to 1L, "piece length" to 1L, + "pieces" to sha1Digest(byteArrayOf(1)) + ))) + ) +} diff --git a/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/ConformanceClients.kt b/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/ConformanceClients.kt new file mode 100644 index 000000000..55ec4d4f7 --- /dev/null +++ b/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/ConformanceClients.kt @@ -0,0 +1,16 @@ +package com.linroid.ketch.torrent + +import java.util.Properties + +/** Shared with the fixture installer and Gradle's test-only dependencies. */ +internal object ConformanceClients { + private val pins = Properties().apply { + val resource = checkNotNull(ConformanceClients::class.java + .getResourceAsStream("/clients.properties")) { + "Missing pinned conformance client versions" + } + resource.use { load(it) } + } + + fun version(client: String): String = checkNotNull(pins.getProperty("$client.version")) +} diff --git a/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/IndependentSeederTest.kt b/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/IndependentSeederTest.kt index dfd8557d7..09d8c1c10 100644 --- a/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/IndependentSeederTest.kt +++ b/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/IndependentSeederTest.kt @@ -7,6 +7,7 @@ import kotlinx.coroutines.test.runTest import kotlinx.coroutines.withContext import kotlinx.coroutines.withTimeout import okio.Path.Companion.toPath +import org.libtorrent4j.LibTorrent import org.libtorrent4j.SessionManager import org.libtorrent4j.SessionParams import org.libtorrent4j.SettingsPack @@ -25,6 +26,9 @@ class IndependentSeederTest { withContext(Dispatchers.Default) { withTimeout(20_000) { NativeLibraryLoader.ensureLoaded() + assertEquals(ConformanceClients.version("libtorrent4j"), LibTorrent.libtorrent4jVersion()) + println("CONFORMANCE_CLIENT libtorrent ${LibTorrent.version()} " + + "libtorrent4j ${LibTorrent.libtorrent4jVersion()}") val root = Files.createTempDirectory("ketch-interop").toFile() val seed = root.resolve("seed").apply { mkdirs() } val payload = ByteArray(80_037) { (it * 31 + 17).toByte() } diff --git a/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/TorrentBenchmarkFixture.kt b/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/TorrentBenchmarkFixture.kt new file mode 100644 index 000000000..07f473726 --- /dev/null +++ b/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/TorrentBenchmarkFixture.kt @@ -0,0 +1,51 @@ +package com.linroid.ketch.torrent + +import java.io.File +import java.security.MessageDigest +import java.util.Random + +/** Streaming fixture construction; even the 10 GiB fixture uses one piece-sized payload buffer. */ +internal data class TorrentBenchmarkFixture(val metainfo: ByteArray, val sha256: ByteArray) { + companion object { + const val PIECE_BYTES: Int = 256 * 1024 + + fun create(file: File, bytes: Long): TorrentBenchmarkFixture { + require(bytes in 1..10L * 1024 * 1024 * 1024) + val random = Random(162) + val content = MessageDigest.getInstance("SHA-256") + val piece = MessageDigest.getInstance("SHA-1") + val hashes = java.io.ByteArrayOutputStream() + val buffer = ByteArray(PIECE_BYTES) + file.outputStream().buffered().use { output -> + var remaining = bytes + while (remaining > 0) { + random.nextBytes(buffer) + val count = minOf(remaining, buffer.size.toLong()).toInt() + output.write(buffer, 0, count) + content.update(buffer, 0, count) + piece.update(buffer, 0, count) + hashes.write(piece.digest()) + remaining -= count + } + } + val metainfo = Bencode.encode(mapOf("info" to mapOf( + "name" to "payload", "length" to bytes, "piece length" to PIECE_BYTES.toLong(), + "pieces" to hashes.toByteArray() + ))) + return TorrentBenchmarkFixture(metainfo, content.digest()) + } + + fun digest(file: File): ByteArray { + val digest = MessageDigest.getInstance("SHA-256") + val buffer = ByteArray(PIECE_BYTES) + file.inputStream().buffered().use { input -> + while (true) { + val count = input.read(buffer) + if (count == -1) break + digest.update(buffer, 0, count) + } + } + return digest.digest() + } + } +} diff --git a/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/TorrentBenchmarkFixtureTest.kt b/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/TorrentBenchmarkFixtureTest.kt new file mode 100644 index 000000000..90b8a1c59 --- /dev/null +++ b/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/TorrentBenchmarkFixtureTest.kt @@ -0,0 +1,32 @@ +package com.linroid.ketch.torrent + +import java.nio.file.Files +import java.security.MessageDigest +import kotlin.test.Test +import kotlin.test.assertContentEquals +import kotlin.test.assertEquals +import kotlin.test.assertFalse + +class TorrentBenchmarkFixtureTest { + @Test + fun streamingFixtureIncludesTheShortFinalPieceAndDistinctPieceContent() { + val root = Files.createTempDirectory("benchmark-fixture").toFile() + try { + val file = root.resolve("payload") + val bytes = TorrentBenchmarkFixture.PIECE_BYTES.toLong() * 2 + 37 + val fixture = TorrentBenchmarkFixture.create(file, bytes) + val metadata = TorrentMetadata.fromBencode(fixture.metainfo) + assertEquals(bytes, file.length()) + assertEquals(bytes, metadata.totalBytes) + assertContentEquals(fixture.sha256, TorrentBenchmarkFixture.digest(file)) + val payload = file.readBytes() // This small test alone reads its bounded fixture into memory. + val hashes = payload.asList().chunked(TorrentBenchmarkFixture.PIECE_BYTES).flatMap { + MessageDigest.getInstance("SHA-1").digest(it.toByteArray()).asList() + }.toByteArray() + assertContentEquals(hashes, metadata.pieceHashes) + assertFalse(hashes.copyOfRange(0, 20).contentEquals(hashes.copyOfRange(20, 40))) + assertContentEquals(fixture.sha256, + TorrentBenchmarkFixture.create(root.resolve("again"), bytes).sha256) + } finally { root.deleteRecursively() } + } +} diff --git a/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/TorrentBenchmarkSeeder.kt b/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/TorrentBenchmarkSeeder.kt new file mode 100644 index 000000000..285723975 --- /dev/null +++ b/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/TorrentBenchmarkSeeder.kt @@ -0,0 +1,118 @@ +package com.linroid.ketch.torrent + +import kotlinx.coroutines.delay +import kotlinx.coroutines.withTimeout +import kotlinx.serialization.json.Json +import kotlinx.serialization.json.JsonObject +import kotlinx.serialization.json.buildJsonObject +import kotlinx.serialization.json.jsonArray +import kotlinx.serialization.json.jsonObject +import kotlinx.serialization.json.jsonPrimitive +import kotlinx.serialization.json.put +import java.io.File +import java.net.ServerSocket +import java.net.Inet4Address +import java.net.NetworkInterface +import java.net.URI +import java.net.http.HttpClient +import java.net.http.HttpRequest +import java.net.http.HttpResponse +import java.time.Duration +import java.util.concurrent.TimeUnit + +/** Out-of-process reference seeder: no native torrent library is loaded into the measuring JVM. */ +internal class TorrentBenchmarkSeeder( + private val root: File, + val dataHost: String = benchmarkAddress(), +) { + private val rpcPort = ServerSocket(0).use { it.localPort } + val peerPort: Int = ServerSocket(0).use { it.localPort } + private val client = HttpClient.newBuilder().connectTimeout(Duration.ofSeconds(2)).build() + private var token: String? = null + private var process: Process? = null + private val log = root.resolve("transmission.log") + + val pid: Long get() = checkNotNull(process).pid() + val cpuNanos: Long get() = checkNotNull(process).info().totalCpuDuration() + .map { it.toNanos() }.orElse(-1L) + + suspend fun start(metainfo: ByteArray?, seed: File) = withTimeout(120_000) { + val binary = checkNotNull(System.getenv("TRANSMISSION_DAEMON")?.takeIf { it.isNotBlank() }) + val version = ProcessBuilder(binary, "--version").redirectErrorStream(true).start() + if (!version.waitFor(5, TimeUnit.SECONDS)) { + version.destroyForcibly() + error("Transmission version check timed out") + } + val text = version.inputStream.bufferedReader().readText() + check(version.exitValue() == 0 && + text.startsWith("transmission-daemon ${ConformanceClients.version("transmission")} ")) + process = ProcessBuilder(binary, "--foreground", "--config-dir", + root.resolve("config").absolutePath, "--download-dir", seed.absolutePath, + "--port", rpcPort.toString(), "--peerport", peerPort.toString(), + "--rpc-bind-address", "127.0.0.1", "--bind-address-ipv4", dataHost, + "--bind-address-ipv6", "::1", + "--no-auth", "--no-dht", "--no-lpd", "--no-portmap", "--no-utp", + "--encryption-tolerated", "--no-global-seedratio") + .redirectErrorStream(true).redirectOutput(log).start() + while (true) { + check(checkNotNull(process).isAlive) { log.readText() } + try { rpc("session-get"); break } catch (_: java.io.IOException) { delay(50) } + } + if (metainfo != null) { + addTorrent(metainfo, seed) + val size = TorrentMetadata.fromBencode(metainfo).totalBytes + while (verifiedBytes() != size) delay(50) + } + } + + fun addTorrent(metainfo: ByteArray, output: File) { + rpc("torrent-add", buildJsonObject { + put("metainfo", encodeBase64(metainfo)); put("download-dir", output.absolutePath) + }) + } + + fun verifiedBytes(): Long { + check(checkNotNull(process).isAlive) { log.readText() } + val torrents = rpc("torrent-get", Json.parseToJsonElement( + "{\"fields\":[\"haveValid\"]}").jsonObject)["torrents"]!!.jsonArray + return torrents.firstOrNull()?.jsonObject?.get("haveValid")?.jsonPrimitive + ?.content?.toLong() ?: 0 + } + + fun close() { + process?.let { + it.destroy() + if (!it.waitFor(5, TimeUnit.SECONDS)) { + it.destroyForcibly() + check(it.waitFor(5, TimeUnit.SECONDS)) { "Transmission did not terminate" } + } + } + } + + private fun rpc(method: String, arguments: JsonObject = buildJsonObject {}): JsonObject { + val body = buildJsonObject { put("method", method); put("arguments", arguments) }.toString() + repeat(2) { + val request = HttpRequest.newBuilder(URI("http://127.0.0.1:$rpcPort/transmission/rpc")) + .timeout(Duration.ofSeconds(5)).POST(HttpRequest.BodyPublishers.ofString(body)) + token?.let { request.header("X-Transmission-Session-Id", it) } + val response = client.send(request.build(), HttpResponse.BodyHandlers.ofString()) + if (response.statusCode() == 409) { + token = response.headers().firstValue("X-Transmission-Session-Id").orElseThrow() + } else { + check(response.statusCode() == 200) + val value = Json.parseToJsonElement(response.body()).jsonObject + check(value["result"]?.jsonPrimitive?.content == "success") + return value["arguments"]!!.jsonObject + } + } + error("Transmission RPC session negotiation failed") + } +} + +/** Transmission rejects loopback tracker peers; select an address actually assigned to this host. */ +internal fun benchmarkAddress(): String = NetworkInterface.networkInterfaces().use { interfaces -> + interfaces.filter { it.isUp && !it.isLoopback }.flatMap { it.inetAddresses.asIterator().asSequence() + .filter { address -> address is Inet4Address && address.isSiteLocalAddress }.toList().stream() + }.findFirst().orElseThrow { IllegalStateException("Benchmark needs a private IPv4 host interface") } + .hostAddress +} diff --git a/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/TorrentBenchmarkTest.kt b/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/TorrentBenchmarkTest.kt index d167245d3..a8c16a8e5 100644 --- a/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/TorrentBenchmarkTest.kt +++ b/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/TorrentBenchmarkTest.kt @@ -3,136 +3,282 @@ package com.linroid.ketch.torrent import com.sun.management.UnixOperatingSystemMXBean import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.delay +import kotlinx.coroutines.flow.first +import kotlinx.coroutines.flow.collect import kotlinx.coroutines.launch import kotlinx.coroutines.runBlocking import kotlinx.coroutines.test.runTest import kotlinx.coroutines.withContext import kotlinx.coroutines.withTimeout -import org.libtorrent4j.SessionManager -import org.libtorrent4j.SessionParams -import org.libtorrent4j.SettingsPack -import org.libtorrent4j.TcpEndpoint -import org.libtorrent4j.TorrentInfo -import org.libtorrent4j.swig.settings_pack +import kotlinx.serialization.Serializable +import kotlinx.serialization.json.Json import java.io.File import java.lang.management.ManagementFactory import java.net.URLClassLoader import java.nio.file.Files import java.util.concurrent.TimeUnit import kotlin.test.Test +import kotlin.test.assertContentEquals import kotlin.test.assertEquals -import kotlin.test.assertTrue +import kotlin.time.Duration.Companion.minutes import kotlin.time.TimeSource -/** Opt-in comparison with isolated downloader processes and a shared reference seed. */ +/** Opt-in, isolated downloader processes with streamed fixtures and a pinned, shared TCP seeder. */ class TorrentBenchmarkTest { @Test - fun compareKotlinAndNativeDownloaders() = runTest { - if (System.getenv("KETCH_TORRENT_BENCHMARK") != "1") return@runTest + fun compareKotlinAndNativeDownloaders() = runTest(timeout = 180.minutes) { + check(System.getenv("KETCH_TORRENT_BENCHMARK") == "1") + val reportFile = File(System.getenv("KETCH_BENCHMARK_REPORT")?.takeIf { it.isNotBlank() } ?: + "build/reports/torrent-benchmark.json").absoluteFile + reportFile.parentFile.mkdirs() + reportFile.delete() // Invalid input must not leave an earlier complete report available. + val sizes = (System.getenv("KETCH_BENCHMARK_BYTES")?.takeIf { it.isNotBlank() } ?: "8388645") + .split(',').map(String::toLong) + val runs = (System.getenv("KETCH_BENCHMARK_RUNS")?.takeIf { it.isNotBlank() } ?: "1").toInt() + require(sizes.isNotEmpty() && sizes.size <= 8 && sizes.distinct().size == sizes.size) + require(sizes.all { it in 1..10L * 1024 * 1024 * 1024 } && runs in 1..10) + val revision = checkNotNull(System.getenv("KETCH_BENCHMARK_REVISION")) + require(revision.matches(Regex("[0-9a-f]{40}"))) withContext(Dispatchers.IO) { - NativeLibraryLoader.ensureLoaded() - val root = Files.createTempDirectory("ketch-benchmark").toFile() - val payload = ByteArray(8 * 1024 * 1024 + 37) { (it * 31 + 17).toByte() } - val seed = root.resolve("seed").apply { mkdirs() } - seed.resolve("payload").writeBytes(payload) - val hashes = payload.asList().chunked(256 * 1024).fold(ByteArray(0)) { value, chunk -> - value + sha1Digest(chunk.toByteArray()) + val results = mutableListOf() + val report = BenchmarkReport(revision, System.getProperty("os.name"), + System.getProperty("os.version"), System.getProperty("os.arch"), + System.getProperty("java.runtime.version"), Runtime.getRuntime().availableProcessors(), + "transmission", ConformanceClients.version("transmission"), sizes, runs) + var failureDetails: String? = null + fun save(complete: Boolean) { + reportFile.writeText(Json.encodeToString(report.copy( + complete = complete, results = results, failure = failureDetails, + ))) } - val data = Bencode.encode(mapOf("info" to mapOf("name" to "payload", - "length" to payload.size.toLong(), "piece length" to 256 * 1024L, "pieces" to hashes))) - root.resolve("fixture.torrent").writeBytes(data) - val manager = SessionManager() - try { - manager.start(SessionParams(referenceSettings())) - val info = TorrentInfo(data) - manager.download(info, seed) - withTimeout(20_000) { - while (manager.find(info.infoHash())?.status()?.isSeeding() != true) delay(20) - } - val urls = generateSequence(javaClass.classLoader) { it.parent } - .filterIsInstance().flatMap { it.urLs.asSequence() } - .map { File(it.toURI()).absolutePath }.toList() - val classpath = (urls + System.getProperty("java.class.path").split(File.pathSeparator)) - .distinct().joinToString(File.pathSeparator) - for (mode in listOf("kotlin", "native")) { - val output = root.resolve(mode).apply { mkdirs() } - val log = root.resolve("$mode.log") - val java = File(System.getProperty("java.home"), "bin/java").absolutePath - val process = ProcessBuilder(java, - "-Xmx256m", "-cp", classpath, "com.linroid.ketch.torrent.TorrentBenchmarkProcess", - mode, root.resolve("fixture.torrent").absolutePath, output.absolutePath, - manager.swig().listen_port().toString()) - .redirectErrorStream(true).redirectOutput(log).start() - var peakRssKiB = 0L - while (process.isAlive) { - val ps = ProcessBuilder("ps", "-o", "rss=", "-p", process.pid().toString()).start() - val rss = ps.inputStream.bufferedReader().readText().trim().toLongOrNull() ?: 0 - ps.waitFor(2, TimeUnit.SECONDS) - peakRssKiB = maxOf(peakRssKiB, rss) - delay(50) + save(false) + for (size in sizes) { + val root = Files.createTempDirectory("ketch-benchmark").toFile() + val seeder = TorrentBenchmarkSeeder(root) + val tracker = TorrentBenchmarkTracker(seeder.peerPort, seeder.dataHost) + try { + require(root.usableSpace >= size * 2 + 2L * 1024 * 1024 * 1024) { + "Benchmark needs space for one seed, one download, and 2 GiB headroom" + } + val seed = root.resolve("seed").apply { mkdirs() } + val generation = TimeSource.Monotonic.markNow() + val fixture = TorrentBenchmarkFixture.create(seed.resolve("payload"), size) + val generationMs = generation.elapsedNow().inWholeMilliseconds + val fixtureHash = fixture.sha256.joinToString("") { + (it.toInt() and 255).toString(16).padStart(2, '0') } - assertEquals(0, process.exitValue(), log.readText()) - assertTrue(output.resolve("payload").readBytes().contentEquals(payload)) - println("BENCHMARK $mode rss_peak_kib=$peakRssKiB ${log.readText().trim()}") + val metadata = TorrentMetadata.fromBencode(fixture.metainfo) + root.resolve("fixture.torrent").writeBytes( + metainfoFromInfo(metadata.infoBytes, listOf(listOf(tracker.url)))) + val seeding = TimeSource.Monotonic.markNow() + seeder.start(fixture.metainfo, seed) + val seedMs = seeding.elapsedNow().inWholeMilliseconds + val urls = generateSequence(javaClass.classLoader) { it.parent } + .filterIsInstance().flatMap { it.urLs.asSequence() } + .map { File(it.toURI()).absolutePath }.toList() + val classpath = (urls + System.getProperty("java.class.path").split(File.pathSeparator)) + .distinct().joinToString(File.pathSeparator) + for (run in 1..runs) { + val modes = if (run % 2 == 1) listOf("kotlin", "native") else listOf("native", "kotlin") + for (mode in modes) { + val output = root.resolve("download").apply { mkdirs() } + val log = root.resolve("$mode-$run.log") + val java = File(System.getProperty("java.home"), "bin/java").absolutePath + val process = ProcessBuilder(java, "-Xmx256m", "-cp", classpath, + "com.linroid.ketch.torrent.TorrentBenchmarkProcess", mode, + root.resolve("fixture.torrent").absolutePath, output.absolutePath, + seeder.dataHost, size.toString()) + .redirectErrorStream(true).redirectOutput(log).start() + var idleRss = -1L + var peakRss = 0L + var measuredProcess: ProcessHandle? = null + try { + withTimeout(630_000) { + while (process.isAlive) { + if (measuredProcess == null) { + val measuredPid = log.readLines().firstOrNull { + it.startsWith("BENCHMARK_READY ") + }?.substringAfter(' ')?.toLong() + if (measuredPid == null) { delay(100); continue } + measuredProcess = if (measuredPid == process.pid()) process.toHandle() else + process.toHandle().descendants().use { children -> + children.filter { it.pid() == measuredPid }.findFirst().orElseThrow() + } + } + val observed = checkNotNull(measuredProcess) + if (!observed.isAlive) { delay(100); continue } + val ps = ProcessBuilder("ps", "-o", "rss=", "-p", observed.pid().toString()) + .start() + if (!ps.waitFor(2, TimeUnit.SECONDS)) { + ps.destroyForcibly() + error("RSS sampler timed out") + } + val rss = ps.inputStream.bufferedReader().readText().trim().toLongOrNull() + if (rss != null) { + peakRss = maxOf(peakRss, rss) + if (idleRss < 0) { + idleRss = rss + process.outputStream.write(10) + process.outputStream.flush() + } + } + delay(100) + } + } + assertEquals(0, process.exitValue(), log.readText()) + check(idleRss >= 0) { "Missing initialized idle RSS sample" } + val metrics = log.readLines().single { it.startsWith("BENCHMARK_METRICS ") } + .removePrefix("BENCHMARK_METRICS ") + val measurement = Json.decodeFromString(metrics) + val verification = TimeSource.Monotonic.markNow() + assertEquals(size, output.resolve("payload").length()) + assertContentEquals(fixture.sha256, + TorrentBenchmarkFixture.digest(output.resolve("payload"))) + results += BenchmarkRun(size, run, mode, fixtureHash, generationMs, seedMs, + verification.elapsedNow().inWholeMilliseconds, idleRss, peakRss, measurement) + save(false) + println("BENCHMARK ${Json.encodeToString(results.last())}") + } catch (failure: Throwable) { + failureDetails = "$mode run $run bytes $size: $failure\n" + + log.readText().takeLast(16_000) + save(false) + throw failure + } finally { + if (process.isAlive) { + process.toHandle().descendants().forEach { child -> + child.destroyForcibly() + child.onExit().get(5, TimeUnit.SECONDS) + } + process.destroyForcibly() + check(process.waitFor(5, TimeUnit.SECONDS)) { "Downloader did not terminate" } + } + check(output.deleteRecursively()) { "Cannot remove completed benchmark payload" } + } + } + } + } catch (failure: Throwable) { + if (failureDetails == null) failureDetails = "Fixture bytes $size: $failure" + save(false) + throw failure + } finally { + seeder.close() + tracker.close() + root.deleteRecursively() } - } finally { manager.stop(); root.deleteRecursively() } + } + save(true) } } } -internal fun referenceSettings(): SettingsPack = SettingsPack().apply { - setString(settings_pack.string_types.listen_interfaces.swigValue(), "127.0.0.1:0") - for (flag in listOf(settings_pack.bool_types.enable_dht, settings_pack.bool_types.enable_lsd, - settings_pack.bool_types.enable_upnp, settings_pack.bool_types.enable_natpmp)) { - setBoolean(flag.swigValue(), false) - } -} +@Serializable +internal data class BenchmarkReport( + val revision: String, val os: String, val osVersion: String, val arch: String, + val java: String, val processors: Int, val referenceEngine: String, val transmission: String, + val sizes: List, val runs: Int, val complete: Boolean = false, + val results: List = emptyList(), val failure: String? = null, +) + +@Serializable +internal data class BenchmarkRun( + val bytes: Long, val run: Int, val mode: String, val fixtureSha256: String, + val fixtureMs: Long, val seedReadyMs: Long, + val verifyMs: Long, val idleRssKiB: Long, val peakRssKiB: Long, val metrics: BenchmarkMetrics, +) + +@Serializable +internal data class BenchmarkMetrics( + val initializationMs: Long, val transferMs: Long, val firstVerifiedMs: Long, + val firstVerifiedBytes: Long, val cpuMs: Long, val idleHeapBytes: Long, + val peakHeapBytes: Long, val idleFds: Long, val peakFds: Long, val finalFds: Long, +) internal object TorrentBenchmarkProcess { @JvmStatic fun main(args: Array): Unit = runBlocking { - withTimeout(90_000) { + withTimeout(600_000) { + val initialization = TimeSource.Monotonic.markNow() val data = File(args[1]).readBytes() val output = File(args[2]) - val port = args[3].toInt() + val payloadBytes = args[4].toLong() + val os = ManagementFactory.getOperatingSystemMXBean() as? UnixOperatingSystemMXBean val memory = ManagementFactory.getMemoryMXBean() - var peakHeap = 0L - var peakFds = 0L - val monitor = launch(Dispatchers.Default) { - while (true) { - peakHeap = maxOf(peakHeap, memory.heapMemoryUsage.used) - peakFds = maxOf(peakFds, os?.openFileDescriptorCount ?: 0) - delay(20) + val native = if (args[0] == "native") { + TorrentBenchmarkSeeder(output.resolve("control").apply { mkdirs() }, args[3]).also { + it.start(null, output) } - } - val start = TimeSource.Monotonic.markNow() + } else null + val kotlin = if (native == null) KotlinTorrentEngine(TorrentConfig(dhtEnabled = false, + connectionsPerTorrent = 1, maxBufferedBytes = 64 * 1024 * 1024, + maxExchangeBytes = 256 * 1024 * 1024, maxSessionStateBytes = 64 * 1024 * 1024), + allowLocalDiscovery = true) else null try { - if (args[0] == "native") { - NativeLibraryLoader.ensureLoaded() - val manager = SessionManager() - try { - manager.start(SessionParams(referenceSettings())) - val info = TorrentInfo(data) - manager.download(info, output) - while (manager.find(info.infoHash()) == null) delay(10) - val handle = manager.find(info.infoHash())!! - handle.swig().connect_peer(TcpEndpoint("127.0.0.1", port).swig()) - while (!handle.status().isSeeding()) delay(10) - } finally { manager.stop() } - } else { - val source = TorrentDownloadSource(TorrentConfig(dhtEnabled = false)) - try { - val metadata = TorrentMetadata.fromBencode(data) - val magnet = MagnetUri(metadata.infoHash, - explicitPeers = listOf("127.0.0.1:$port")).toUri() - source.download(sourceContext(magnet, source.resolveMetainfo(data), - output.resolve("payload").absolutePath)) - } finally { source.close() } + kotlin?.start() + val initializationMs = initialization.elapsedNow().inWholeMilliseconds + val idleHeap = if (native == null) memory.heapMemoryUsage.used else -1 + val idleFds = if (native == null) os?.openFileDescriptorCount ?: -1 else -1 + println("BENCHMARK_READY ${native?.pid ?: ProcessHandle.current().pid()}") + System.out.flush() + check(System.`in`.read() == 10) { "Missing measurement start signal" } + var peakHeap = idleHeap + var peakFds = idleFds + val monitor = launch(Dispatchers.Default) { + while (true) { + if (native == null) { + peakHeap = maxOf(peakHeap, memory.heapMemoryUsage.used) + peakFds = maxOf(peakFds, os?.openFileDescriptorCount ?: -1) + } + delay(20) + } } - } finally { monitor.cancel(); monitor.join() } - println("elapsed_ms=${start.elapsedNow().inWholeMilliseconds} heap_peak_bytes=$peakHeap " + - "fd_peak=$peakFds fd_after=${os?.openFileDescriptorCount}") + val cpuStart = native?.cpuNanos ?: os?.processCpuTime ?: -1 + val transfer = TimeSource.Monotonic.markNow() + var firstVerifiedMs = -1L + var firstVerifiedBytes = 0L + fun verified(bytes: Long) { + if (bytes > 0 && firstVerifiedMs < 0) { + firstVerifiedMs = transfer.elapsedNow().inWholeMilliseconds + firstVerifiedBytes = bytes + } + } + try { + if (native != null) { + native.addTorrent(data, output) + while (true) { + val bytes = native.verifiedBytes() + verified(bytes) + if (bytes == payloadBytes) break + delay(100) + } + } else { + val metadata = TorrentMetadata.fromBencode(data) + val session = checkNotNull(kotlin).addTask(TorrentTaskSpec("benchmark", metadata, + output.resolve("payload").absolutePath, emptySet())) + val progress = launch { session.downloadedBytes.collect { verified(it) } } + try { + session.resume() + val state = session.state.first { + it == TorrentSessionState.FINISHED || it == TorrentSessionState.STOPPED + } + verified(session.downloadedBytes.value) + check(state == TorrentSessionState.FINISHED) { session.failure.value.toString() } + } finally { progress.cancel(); progress.join() } + } + } finally { monitor.cancel(); monitor.join() } + val elapsed = transfer.elapsedNow().inWholeMilliseconds + val cpuEnd = native?.cpuNanos ?: os?.processCpuTime ?: -1 + val cpuMs = if (cpuStart < 0 || cpuEnd < 0) -1 else (cpuEnd - cpuStart) / 1_000_000 + native?.close() + kotlin?.stop() + val metrics = BenchmarkMetrics(initializationMs, elapsed, firstVerifiedMs, + firstVerifiedBytes, cpuMs, idleHeap, + peakHeap, idleFds, peakFds, if (native == null) os?.openFileDescriptorCount ?: -1 else -1) + println("BENCHMARK_METRICS ${Json.encodeToString(metrics)}") + } finally { + native?.close() + kotlin?.stop() + } } } } diff --git a/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/TorrentBenchmarkTracker.kt b/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/TorrentBenchmarkTracker.kt new file mode 100644 index 000000000..2c8e5bfa9 --- /dev/null +++ b/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/TorrentBenchmarkTracker.kt @@ -0,0 +1,26 @@ +package com.linroid.ketch.torrent + +import com.sun.net.httpserver.HttpServer +import java.net.InetSocketAddress +import java.net.InetAddress + +/** Both downloaders discover the same seeder through this loopback-only fixture tracker. */ +internal class TorrentBenchmarkTracker(peerPort: Int, dataHost: String) { + private val server = HttpServer.create(InetSocketAddress("127.0.0.1", 0), 0) + val url: String get() = "http://127.0.0.1:${server.address.port}/announce" + + init { + val peer = InetAddress.getByName(dataHost).address + + byteArrayOf((peerPort shr 8).toByte(), peerPort.toByte()) + val response = Bencode.encode(mapOf("interval" to 60L, "peers" to peer)) + server.createContext("/announce") { exchange -> + try { + exchange.sendResponseHeaders(200, response.size.toLong()) + exchange.responseBody.use { it.write(response) } + } finally { exchange.close() } + } + server.start() + } + + fun close() { server.stop(0) } +} diff --git a/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/TorrentStorageCancellationTest.kt b/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/TorrentStorageCancellationTest.kt new file mode 100644 index 000000000..10c61dfd9 --- /dev/null +++ b/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/TorrentStorageCancellationTest.kt @@ -0,0 +1,61 @@ +package com.linroid.ketch.torrent + +import kotlinx.coroutines.CompletableDeferred +import kotlinx.coroutines.CoroutineStart +import kotlinx.coroutines.async +import kotlinx.coroutines.cancelAndJoin +import kotlinx.coroutines.sync.Semaphore +import kotlinx.coroutines.test.runTest +import okio.FileHandle +import okio.FileSystem +import okio.ForwardingFileSystem +import okio.Path +import java.util.concurrent.CountDownLatch +import java.util.concurrent.TimeUnit +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertFalse +import kotlin.test.assertTrue + +class TorrentStorageCancellationTest { + @Test + fun canceledBlockedProviderRetainsSlotUntilOperationReturns() = runTest { + val entered = CompletableDeferred() + val unblock = CountDownLatch(1) + val slots = Semaphore(1) + val root = FileSystem.SYSTEM_TEMPORARY_DIRECTORY / + "ketch-storage-blocked-${InfoHash.fromBytes(torrentRandomBytes(20)).hex}" + val metadata = TorrentMetadata.fromBencode(Bencode.encode(mapOf("info" to mapOf( + "name" to "empty", "piece length" to 16_384L, "length" to 0L, "pieces" to byteArrayOf() + )))) + val provider = object : ForwardingFileSystem(torrentFileSystem) { + override fun openReadWrite(file: Path, mustCreate: Boolean, mustExist: Boolean): FileHandle { + entered.complete(Unit) + check(unblock.await(10, TimeUnit.SECONDS)) { "Test did not release blocked provider" } + return super.openReadWrite(file, mustCreate, mustExist) + } + } + val first = TorrentPieceStore(metadata, root / "first", emptySet(), "first", provider, slots) + val second = TorrentPieceStore(metadata, root / "second", emptySet(), "second", + storageSlots = slots) + val blocked = async { first.initialize() } + try { + entered.await() + blocked.cancel() + assertFalse(blocked.isCompleted) + assertEquals(0, slots.availablePermits) + val waiting = async(start = CoroutineStart.UNDISPATCHED) { second.initialize() } + assertFalse(waiting.isCompleted) + assertFalse(torrentFileSystem.exists(root / "second")) + unblock.countDown() + blocked.join() + waiting.await() + assertTrue(second.isInitialized()) + assertEquals(1, slots.availablePermits) + } finally { + unblock.countDown() + blocked.cancelAndJoin() + torrentFileSystem.deleteRecursively(root, mustExist = false) + } + } +} diff --git a/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/TransmissionInteropTest.kt b/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/TransmissionInteropTest.kt index 1749f2f3f..c85b4ad0a 100644 --- a/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/TransmissionInteropTest.kt +++ b/library/torrent/src/jvmTest/kotlin/com/linroid/ketch/torrent/TransmissionInteropTest.kt @@ -30,13 +30,20 @@ import kotlin.test.assertTrue class TransmissionInteropTest { @Test fun publicSource_resolvesTrackerlessMagnetAndDownloadsFromTransmission() = runTest { - val binary = System.getenv("TRANSMISSION_DAEMON")?.takeIf { it.isNotBlank() } ?: return@runTest + val binary = checkNotNull(System.getenv("TRANSMISSION_DAEMON")?.takeIf { it.isNotBlank() }) { + "Transmission conformance requires TRANSMISSION_DAEMON; see test-fixtures/torrent/README.md" + } withContext(Dispatchers.IO) { withTimeout(60_000) { val version = ProcessBuilder(binary, "--version").redirectErrorStream(true).start() - val versionText = version.inputStream.bufferedReader().readText() - assertTrue(version.waitFor(5, TimeUnit.SECONDS)) - assertTrue("4.1.3" in versionText || "4.0.5" in versionText, versionText) + val exited = version.waitFor(5, TimeUnit.SECONDS) + if (!exited) version.destroyForcibly().waitFor() + assertTrue(exited, "Transmission version check timed out") + val versionText = version.inputStream.bufferedReader().readText().trim() + assertEquals(0, version.exitValue(), versionText) + val expected = ConformanceClients.version("transmission") + assertTrue(versionText.startsWith("transmission-daemon $expected "), versionText) + println("CONFORMANCE_CLIENT $versionText") val root = Files.createTempDirectory("ketch-transmission").toFile() val seed = root.resolve("seed").apply { mkdirs() } val payload = ByteArray(512 * 1024 + 37) { (it * 31 + 17).toByte() } diff --git a/test-fixtures/torrent/BENCHMARK.md b/test-fixtures/torrent/BENCHMARK.md new file mode 100644 index 000000000..083688bc0 --- /dev/null +++ b/test-fixtures/torrent/BENCHMARK.md @@ -0,0 +1,77 @@ +# Direct TCP baseline + +The benchmark compares an isolated Kotlin JVM with the pinned native Transmission downloader, +using a separate pinned Transmission seeder. A loopback-only tracker advertises a private IPv4 +address assigned to this host. Transmission rejects loopback tracker peers, so data sockets bind +to that host address; all data still travels between processes on the same machine. uTP and public +discovery are disabled. Payloads are deterministic +pseudorandom data from Java Random seed 162, generated one 256 KiB piece at a time. The independent +fixture builder uses JDK SHA-1 and SHA-256. Output verification streams SHA-256 and checks file length. +Neither generation nor verification allocates a payload-sized byte array. + +Run on an otherwise idle Unix host with `ps`, JDK 21, and sufficient disk space. Only one output is +retained at a time. Each size needs twice the payload size plus 2 GiB free space. Reference hardware, +OS, runtime, filesystem, and competing workloads must be recorded alongside the JSON results. + +```shell +TRANSMISSION_DAEMON=/tmp/ketch-transmission/build/daemon/transmission-daemon \ +KETCH_TORRENT_BENCHMARK=1 \ +KETCH_BENCHMARK_BYTES=1073741824,10737418240 \ +KETCH_BENCHMARK_RUNS=5 \ +KETCH_BENCHMARK_REVISION="$(git rev-parse HEAD)" \ +KETCH_BENCHMARK_REPORT=/tmp/ketch-torrent-baseline.json \ + ./gradlew :library:torrent:jvmTest --tests '*TorrentBenchmarkTest*' +``` + +For a harness smoke test, omit sizes/runs: defaults are 8 MiB + 37 bytes and one run. Every benchmark +invocation executes; Gradle test caching is disabled. Alternating Kotlin/native order reduces a +systematic order advantage. File data is not explicitly purged from the OS cache between runs. +Report all raw samples and the median/spread, including failures; do not select only fast samples. + +The seeder has a 120-second readiness deadline. Each child has a 600-second execution deadline and +the parent enforces a 630-second wall-clock deadline including readiness and sampling. A stuck +child is forcibly terminated, with a five-second termination deadline. These bounds are fixed +before measurement. Only `complete: true` proves completion. Failures preserve previously completed +samples. +A complete 1/10 GiB baseline requires twenty verified child runs, not just a successful smoke test. + +The Kotlin JVM has a 256 MiB heap ceiling. The native engine is the Transmission executable; its +small JVM RPC controller is excluded from native RSS/CPU measurements. Unsupported native heap/FD +and CPU counters are reported as -1, never as zero. The Kotlin engine uses one peer, disabled uploads +and DHT, 64 MiB transfer capacity, 64 MiB session capacity, and a 256 MiB combined admission ceiling; +other limits retain their defaults. The native reference disables DHT, LPD, NAT mapping and uTP, +uses the same fixture tracker, and otherwise retains Transmission defaults. One local seeder is +available to either downloader, with no upload recipients. +These are explicit benchmark settings, not a completed production resource profile. + +An initialized controller reports the measured engine PID and waits for the parent to sample idle +RSS before transfer. The parent verifies that the PID belongs to the controller or its descendants. +Forced timeout cleanup also terminates controller descendants. +The report includes runtime/platform identity, fixture SHA-256, fixture creation time, seeder readiness, +child initialization, download-to-verified-completion time, CPU time, heap/RSS samples, descriptor +counts, and independent output verification time. RSS is sampled every 100 ms; heap/FDs every 20 ms, +so short peaks may be missed. Idle heap is not a forced-GC measurement. CPU can exceed wall time +because work uses multiple threads. Transfer time includes overlapping hashing and disk work and +is not a measurement of network-only throughput. Final descriptor counts follow engine shutdown. + +This baseline does not certify the 80% sustained-throughput release target. Matching native cache +and durability policy, separating overlapping hash/disk/network costs, shaped swarms, many files, +partial selection, mobile hardware, energy, lifecycle, and soak measurements remain release work. +No target is relaxed to accommodate a slow run. + +The first verified-progress timestamp separates connection/unchoke wait from subsequent transfer. +Native verified bytes are polled through RPC every 100 ms and may also be delayed by Transmission's +own statistics refresh. The 1/10 GiB fixtures contain only full pieces. Report both end-to-end and +post-first-piece rates, without calling either one a complete release benchmark. + +An initial run at `560ba380` used an in-process libtorrent seeder and its parent JVM crashed with +SIGSEGV on the Java Finalizer thread after one verified Kotlin 1 GiB sample. It is incomplete and +excluded from aggregate comparisons. The crash does not establish a root cause in the downloader. +The next attempt moved the seeder to an external pinned Transmission process and isolated native +torrent bindings in the reference downloader JVM. No timeout or performance target was raised. + +A second attempt at `d92566d8` isolated the Transmission seeder but the libtorrent4j downloader JVM +also exited with SIGSEGV on its first 1 GiB sample. That incomplete report is retained separately. +The current native baseline uses the pinned Transmission executable directly, with no torrent JNI +binding in any measurement JVM. Existing libtorrent4j interoperability tests remain separate; the +large-fixture binding crash is not a successful compatibility or performance result. diff --git a/test-fixtures/torrent/README.md b/test-fixtures/torrent/README.md new file mode 100644 index 000000000..bff14ee0b --- /dev/null +++ b/test-fixtures/torrent/README.md @@ -0,0 +1,40 @@ +# Torrent conformance fixtures + +These clients are test fixtures. They must never enter a product runtime dependency graph. +`clients.properties` pins their versions; Transmission's release archive is authenticated by SHA-256. +The libtorrent fixture remains a version-pinned JVM test dependency and validates its loaded version. + +Build Transmission with Python 3.12+, CMake, a C/C++ compiler, and platform curl/TLS development +libraries. On Ubuntu 24.04, install `cmake`, `build-essential`, `libcurl4-openssl-dev`, and +`libssl-dev`. Third-party torrent libraries are built from the verified Transmission source archive. + +```shell +python3 tools/torrent/build_transmission.py --output /tmp/ketch-transmission +TRANSMISSION_DAEMON=/tmp/ketch-transmission/build/daemon/transmission-daemon \ + ./gradlew :library:torrent:jvmTest -PtorrentConformance=true +python3 tools/torrent/verify_conformance.py \ + --reports library/torrent/build/test-results/jvmTest \ + --revision "$(git rev-parse HEAD)" \ + --output library/torrent/build/reports/conformance/executed.json +``` + +`torrentConformance=true` always executes the torrent JVM suite; its results cannot be restored +from Gradle's test cache. A missing Transmission binary or a version different from the pin fails +the required suite. Ordinary JVM runs exclude Transmission when it is not configured, rather than +counting a no-op test as a successful interoperability scenario. Supplying the binary opts in. +`KETCH_TORRENT_BENCHMARK=1` similarly includes the expensive process comparison; otherwise it is +excluded from the test plan. Packaging fixtures remain opt-in until their release lane is added. + +`scenarios.json` lists actual required behavioral tests. The evidence verifier rejects missing, +skipped, failed, and duplicate results and records the Git revision, report digests, client versions, +and durations. It removes an older evidence file before validating, so failure cannot leave a stale +green report available for upload. Run the gate's negative cases with: + +```shell +python3 -m unittest discover -s tools/torrent -p 'test_*.py' +``` + +The first manifest covers **v1 download** through libtorrent DHT/metadata and a Transmission magnet +with an explicit peer. It does not certify v2, hybrid, uploads, or public trackerless discovery. +Extend the manifest with each implemented capability and preserve independent expected payloads. +Never add an unimplemented scenario as passing, or satisfy two capability rows with one no-op test. diff --git a/test-fixtures/torrent/clients.properties b/test-fixtures/torrent/clients.properties new file mode 100644 index 000000000..9ac7a85eb --- /dev/null +++ b/test-fixtures/torrent/clients.properties @@ -0,0 +1,5 @@ +# Test-only clients. Update pins and recorded interoperability evidence together. +transmission.version=4.1.3 +transmission.url=https://github.com/transmission/transmission/releases/download/4.1.3/transmission-4.1.3.tar.xz +transmission.sha256=ce7d2d8b101f7eb54bc3cf0bc55f52f7ebd4a25fa48e00bdca9a7e0fc02617da +libtorrent4j.version=2.1.0-39 diff --git a/test-fixtures/torrent/scenarios.json b/test-fixtures/torrent/scenarios.json new file mode 100644 index 000000000..84a3fe896 --- /dev/null +++ b/test-fixtures/torrent/scenarios.json @@ -0,0 +1,15 @@ +{ + "schemaVersion": 1, + "scenarios": [ + { + "id": "v1.libtorrent.dht-metadata-download", + "class": "com.linroid.ketch.torrent.IndependentSeederTest", + "method": "kotlinDownloader_completesVerifiedTransferFromIndependentSeeder" + }, + { + "id": "v1.transmission.explicit-peer-magnet-download", + "class": "com.linroid.ketch.torrent.TransmissionInteropTest", + "method": "publicSource_resolvesTrackerlessMagnetAndDownloadsFromTransmission" + } + ] +} diff --git a/tools/torrent/.gitignore b/tools/torrent/.gitignore new file mode 100644 index 000000000..c18dd8d83 --- /dev/null +++ b/tools/torrent/.gitignore @@ -0,0 +1 @@ +__pycache__/ diff --git a/tools/torrent/build_transmission.py b/tools/torrent/build_transmission.py new file mode 100644 index 000000000..0150f6596 --- /dev/null +++ b/tools/torrent/build_transmission.py @@ -0,0 +1,65 @@ +"""Build the pinned, test-only Transmission daemon in an isolated output directory.""" + +import argparse +import hashlib +import os +from pathlib import Path +import subprocess +import tarfile +import urllib.request + + +def main(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--output", type=Path, required=True) + args = parser.parse_args() + root = Path(__file__).resolve().parents[2] + pins = dict( + line.split("=", 1) + for line in (root / "test-fixtures/torrent/clients.properties").read_text().splitlines() + if line and not line.startswith("#") + ) + output = args.output.resolve() + output.mkdir(parents=True, exist_ok=True) + archive = output / "source.tar.xz" + if not archive.exists(): + temporary = output / "source.tar.xz.part" + with urllib.request.urlopen(pins["transmission.url"], timeout=60) as response: + with temporary.open("wb") as destination: + while chunk := response.read(1024 * 1024): + destination.write(chunk) + temporary.replace(archive) + with archive.open("rb") as stream: + digest = hashlib.file_digest(stream, "sha256").hexdigest() + if digest != pins["transmission.sha256"]: + raise SystemExit(f"Transmission archive digest mismatch: {digest}") + source = output / f"transmission-{pins['transmission.version']}" + if not source.exists(): + with tarfile.open(archive) as package: + package.extractall(output, filter="data") + build = output / "build" + subprocess.run([ + "cmake", "-S", str(source), "-B", str(build), "-DCMAKE_BUILD_TYPE=Release", + "-DENABLE_DAEMON=ON", "-DENABLE_GTK=OFF", "-DENABLE_QT=OFF", "-DENABLE_MAC=OFF", + "-DENABLE_TESTS=OFF", "-DENABLE_UTILS=OFF", "-DENABLE_CLI=OFF", "-DENABLE_NLS=OFF", + "-DINSTALL_WEB=OFF", "-DINSTALL_DOC=OFF", "-DRUN_CLANG_TIDY=OFF", + # Use the implementations shipped in the authenticated release archive. + *[f"-DUSE_SYSTEM_{name}=OFF" for name in ( + "EVENT2", "DEFLATE", "DHT", "MINIUPNPC", "NATPMP", "UTP", "B64", "PSL" + )], + ], check=True) + subprocess.run([ + "cmake", "--build", str(build), "--target", "transmission-daemon", "--parallel", + str(min(os.cpu_count() or 2, 4)), + ], check=True) + binary = build / "daemon/transmission-daemon" + result = subprocess.run([str(binary), "--version"], check=True, capture_output=True, + text=True, timeout=10) + version = (result.stdout + result.stderr).strip() + if not version.startswith(f"transmission-daemon {pins['transmission.version']} "): + raise SystemExit(f"Unexpected Transmission version: {version}") + print(binary) + + +if __name__ == "__main__": + main() diff --git a/tools/torrent/test_verify_conformance.py b/tools/torrent/test_verify_conformance.py new file mode 100644 index 000000000..f2ea40aaf --- /dev/null +++ b/tools/torrent/test_verify_conformance.py @@ -0,0 +1,59 @@ +"""Behavior checks for the release evidence gate; no external peers are needed.""" + +from pathlib import Path +import tempfile +import unittest + +from verify_conformance import verify + + +class ConformanceGateTest(unittest.TestCase): + def setUp(self): + self.directory = tempfile.TemporaryDirectory() + self.addCleanup(self.directory.cleanup) + self.reports = Path(self.directory.name) + self.manifest = {"schemaVersion": 1, "scenarios": [ + {"id": "download", "class": "PeerTest", "method": "download"} + ]} + + def report(self, content, name="TEST-peer.xml"): + (self.reports / name).write_text(f"{content}") + + def test_missing_report_fails(self): + with self.assertRaisesRegex(ValueError, "did not execute"): + verify(self.manifest, self.reports) + + def test_skipped_required_case_fails(self): + self.report('') + with self.assertRaisesRegex(ValueError, "did not execute"): + verify(self.manifest, self.reports) + + def test_failure_elsewhere_cannot_produce_green_evidence(self): + self.report('' + '') + with self.assertRaisesRegex(ValueError, "Failed tests"): + verify(self.manifest, self.reports) + + def test_duplicate_reports_fail(self): + case = '' + self.report(case) + self.report(case, "TEST-duplicate.xml") + with self.assertRaisesRegex(ValueError, "Duplicate reported"): + verify(self.manifest, self.reports) + + def test_one_case_cannot_satisfy_two_scenarios(self): + self.report('') + self.manifest["scenarios"].append( + {"id": "different", "class": "PeerTest", "method": "download"}) + with self.assertRaisesRegex(ValueError, "Duplicate scenario"): + verify(self.manifest, self.reports) + + def test_executed_target_case_records_duration_and_report_digest(self): + self.report('') + evidence = verify(self.manifest, self.reports) + self.assertEqual([{"id": "download", "seconds": 1.25}], evidence["scenarios"]) + self.assertEqual(64, len(evidence["reports"][0]["sha256"])) + + +if __name__ == "__main__": + unittest.main() diff --git a/tools/torrent/verify_conformance.py b/tools/torrent/verify_conformance.py new file mode 100644 index 000000000..652c5618e --- /dev/null +++ b/tools/torrent/verify_conformance.py @@ -0,0 +1,79 @@ +"""Require successful, executed JUnit scenarios and write revision-bound conformance evidence.""" + +import argparse +from datetime import datetime, timezone +import hashlib +import json +from pathlib import Path +import re +import xml.etree.ElementTree as ET + + +def verify(manifest, reports): + if manifest.get("schemaVersion") != 1 or not manifest.get("scenarios"): + raise ValueError("Unknown or empty conformance manifest") + cases = {} + report_evidence = [] + clients = set() + for report in sorted(reports.glob("TEST-*.xml")): + data = report.read_bytes() + document = ET.fromstring(data) + report_evidence.append({"file": report.name, "sha256": hashlib.sha256(data).hexdigest()}) + for case in document.iter("testcase"): + # Kotlin's JUnit runner appends the target name to the method. + method = re.sub(r"\[[^\]]+\]$", "", case.get("name", "")) + key = (case.get("classname"), method) + if key in cases: + raise ValueError(f"Duplicate reported testcase: {key}") + cases[key] = case + for output in document.iter("system-out"): + clients.update(line.removeprefix("CONFORMANCE_CLIENT ") + for line in (output.text or "").splitlines() + if line.startswith("CONFORMANCE_CLIENT ")) + if (next(document.iter("failure"), None) is not None or + next(document.iter("error"), None) is not None): + raise ValueError(f"Failed tests in {report.name}") + executed = [] + identities = set() + test_keys = set() + for scenario in manifest["scenarios"]: + identity = scenario["id"] + key = (scenario["class"], scenario["method"]) + if identity in identities or key in test_keys: + raise ValueError(f"Duplicate scenario identity or test: {identity}") + identities.add(identity) + test_keys.add(key) + case = cases.get(key) + if case is None or case.find("skipped") is not None: + raise ValueError(f"Required scenario did not execute: {identity}") + executed.append({"id": identity, "seconds": float(case.get("time", "0"))}) + return {"scenarios": executed, "reports": report_evidence, "clients": sorted(clients)} + + +def main(): + parser = argparse.ArgumentParser(description=__doc__) + root = Path(__file__).resolve().parents[2] + parser.add_argument("--manifest", type=Path, + default=root / "test-fixtures/torrent/scenarios.json") + parser.add_argument("--reports", type=Path, required=True) + parser.add_argument("--revision", required=True) + parser.add_argument("--output", type=Path, required=True) + args = parser.parse_args() + # A failed invocation must not leave an earlier green evidence file available for upload. + args.output.unlink(missing_ok=True) + if not re.fullmatch(r"[0-9a-f]{40,64}", args.revision): + raise ValueError("Evidence requires a full Git revision") + raw_manifest = args.manifest.read_bytes() + evidence = verify(json.loads(raw_manifest), args.reports) + evidence.update({ + "schemaVersion": 1, "revision": args.revision, + "recordedAt": datetime.now(timezone.utc).isoformat(), + "manifestSha256": hashlib.sha256(raw_manifest).hexdigest(), + }) + args.output.parent.mkdir(parents=True, exist_ok=True) + args.output.write_text(json.dumps(evidence, indent=2) + "\n") + print(f"Verified {len(evidence['scenarios'])} conformance scenarios at {args.revision}") + + +if __name__ == "__main__": + main()