From 6c07dcb839ea694995733dcea3b915eb92d02f0b Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 1 Sep 2026 17:26:50 +0000 Subject: [PATCH 1/6] perf(quartz): drop copy-on-write from linuxX64's LargeCache MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit LocalCache's linuxX64 store kept a LinkedHashMap inside an AtomicReference and replaced it wholesale on every write, so each put was O(n) in the size of the cache and filling it was O(n^2). It was not thread-safe either: the read-copy-write was not a CAS loop, so concurrent writers silently dropped each other's entries. Replace it with a mutable map guarded by PlatformLock plus a lazily rebuilt read snapshot. Point operations (get/put/remove/containsKey/size) are O(1) under the lock; bulk operations run against a point-in-time copy rebuilt at most once per write epoch, which also keeps caller-supplied lambdas out of the critical section — PlatformLock is not reentrant here and LocalCache predicates call back into the cache. Two behaviour fixes fall out of matching the JVM actual's putIfAbsent: createIfAbsent now reports true only when this call inserted (it previously returned get(key) != null, which also reported true when another thread had just created the entry), and getOrCreate publishes atomically. ConcurrentHashCache.linux gets the same treatment. Its only caller, CachingEventDecoder, writes once per event arriving from a relay, so the per-write map rebuild was the worst-placed copy of the three. None of this was caught because no CI job compiled or ran linuxX64. Add LargeCacheTest to commonTest as a cross-target contract for the ~40 methods each actual reimplements by hand, a linuxTest suite covering the concurrency this actual now has to get right on its own, and a CI leg that runs both on Linux Native. That leg is scoped to the cache and concurrency packages: the full linuxX64Test suite is 3,490 tests with 78 pre-existing failures, nearly all TODO() stubs in linux actuals that were never written (MLS crypto, the SQLite driver, NIP-44, Bolt12). Filling those in is its own project; the filter keeps the job meaningful and green, and widening it later is one line. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01HxQ1QuyzSkR38iFHbREjoS --- .github/workflows/build.yml | 70 ++++ .../quartz/utils/cache/LargeCacheTest.kt | 319 ++++++++++++++++++ .../utils/cache/ConcurrentHashCache.linux.kt | 33 +- .../quartz/utils/cache/LargeCache.linux.kt | 169 ++++++++-- .../utils/cache/LargeCacheConcurrencyTest.kt | 120 +++++++ .../cache/LargeCacheRangeFallbackTest.kt | 73 ++++ .../utils/concurrent/ConcurrentMap.native.kt | 10 +- 7 files changed, 742 insertions(+), 52 deletions(-) create mode 100644 quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCacheTest.kt create mode 100644 quartz/src/linuxTest/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCacheConcurrencyTest.kt create mode 100644 quartz/src/linuxTest/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCacheRangeFallbackTest.kt diff --git a/.github/workflows/build.yml b/.github/workflows/build.yml index 06d0b541d1..c130551bce 100644 --- a/.github/workflows/build.yml +++ b/.github/workflows/build.yml @@ -194,6 +194,76 @@ jobs: name: geode Test Reports path: geode/build/reports + # linuxX64 is the only target whose LargeCache / ConcurrentHashCache actuals are + # hand-written concurrent maps rather than a delegation to a platform concurrent + # collection, and until this job existed nothing ran them: the target was compiled + # by no CI leg at all. That is how a copy-on-write LargeCache with O(n) writes and a + # non-atomic read-copy-write (concurrent writers silently dropped entries) sat in + # the tree unnoticed. + # + # Scoped to the cache/concurrency packages on purpose. The full linuxX64Test suite + # is 3,490 tests with 78 pre-existing failures, essentially all of them `TODO()` + # stubs in linux actuals that were never written — MLS crypto, the SQLite driver, + # NIP-44, Bolt12 — plus two URL-handling divergences. Filling those in is its own + # project; gating PRs on them today would just mean a permanently red job. The + # filter keeps the leg meaningful and green, and widening it is a one-line change + # once the native actuals land. + # + # This still compiles and links the whole module for linuxX64, so a commonMain or + # commonTest source that reaches for a JVM-only API fails here too — on a target + # with no Foundation to fall back on the way Apple has. + test-quartz-linux-native: + needs: lint + runs-on: ubuntu-latest + timeout-minutes: 45 + steps: + - name: Checkout code + uses: actions/checkout@v7 + + - name: Set up JDK 21 + uses: actions/setup-java@v6.0.0 + with: + distribution: 'temurin' + java-version: 21 + + - name: Set up Gradle + uses: gradle/actions/setup-gradle@v6 + with: + cache-read-only: ${{ github.ref != 'refs/heads/main' }} + + # The Kotlin/Native toolchain (compiler distribution + LLVM + the sysroot) lands + # in ~/.konan, which setup-gradle does not cache. Without this the job re-downloads + # well over a gigabyte on every run. Keyed on the version catalog so a Kotlin bump + # re-populates it. + - name: Cache Kotlin/Native toolchain + uses: actions/cache@v4 + with: + path: ~/.konan + key: konan-${{ runner.os }}-${{ hashFiles('gradle/libs.versions.toml') }} + restore-keys: konan-${{ runner.os }}- + + - name: Test Quartz caches on Linux Native + run: | + ./gradlew :quartz:linuxX64Test \ + --tests "com.vitorpamplona.quartz.utils.cache.*" \ + --tests "com.vitorpamplona.quartz.utils.concurrent.*" + + - name: Linux Native Test Report + uses: mikepenz/action-junit-report@a9170d5795813c01ab4901ffb045b52bab4ab09d # v6.5.0 + if: always() + with: + report_paths: 'quartz/build/test-results/linuxX64Test/TEST-*.xml' + annotate_only: true + detailed_summary: true + fail_on_failure: true + + - name: Upload Linux Native Test Reports + uses: actions/upload-artifact@v7 + if: failure() + with: + name: Quartz Linux Native Test Reports + path: quartz/build/reports + test-quartz-ios: # Phase 1 of the iOS support plan # (amethyst/plans/2026-05-24-ios-support.md): keep :quartz green on iOS diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCacheTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCacheTest.kt new file mode 100644 index 0000000000..bc61f7aca0 --- /dev/null +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCacheTest.kt @@ -0,0 +1,319 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.quartz.utils.cache + +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertFalse +import kotlin.test.assertNull +import kotlin.test.assertTrue + +/** + * Cross-target contract for [LargeCache], the store behind `LocalCache`. + * + * There was no shared suite for this class: the JVM/Android actual is covered only + * indirectly through `LocalCache`, and the linuxX64 and Apple actuals — hand-written + * reimplementations of the same ~40 methods — were covered by nothing at all. This + * runs the same assertions against whichever actual the target picked, so a divergence + * shows up as a test failure instead of at runtime on one platform. + * + * Deliberately order-agnostic: iteration order is insertion order on linux, sorted-key + * order on JVM/Android (`ConcurrentSkipListMap`) and hash order on Apple. Anything + * order-sensitive is compared as a set or sorted first. + */ +class LargeCacheTest { + private fun cacheOf(vararg pairs: Pair) = + LargeCache().apply { + pairs.forEach { put(it.first, it.second) } + } + + @Test + fun emptyCache() { + val cache = LargeCache() + + assertEquals(0, cache.size()) + assertTrue(cache.isEmpty()) + assertNull(cache.get("a")) + assertFalse(cache.containsKey("a")) + assertTrue(cache.keys().isEmpty()) + assertTrue(cache.values().toList().isEmpty()) + } + + @Test + fun putGetAndOverwrite() { + val cache = cacheOf("a" to 1, "b" to 2) + + assertEquals(1, cache.get("a")) + assertEquals(2, cache.get("b")) + assertEquals(2, cache.size()) + assertFalse(cache.isEmpty()) + assertTrue(cache.containsKey("a")) + + cache.put("a", 10) + + assertEquals(10, cache.get("a")) + assertEquals(2, cache.size(), "overwriting a key must not grow the cache") + } + + @Test + fun removeReturnsOldValue() { + val cache = cacheOf("a" to 1, "b" to 2) + + assertEquals(1, cache.remove("a")) + assertNull(cache.get("a")) + assertFalse(cache.containsKey("a")) + assertEquals(1, cache.size()) + + assertNull(cache.remove("a"), "removing an absent key returns null") + assertEquals(1, cache.size()) + } + + @Test + fun clearEmptiesEverything() { + val cache = cacheOf("a" to 1, "b" to 2) + + cache.clear() + + assertEquals(0, cache.size()) + assertTrue(cache.isEmpty()) + assertNull(cache.get("a")) + assertTrue(cache.keys().isEmpty()) + assertTrue(cache.values().toList().isEmpty()) + } + + @Test + fun getOrCreateBuildsOnlyOnce() { + val cache = LargeCache() + var builds = 0 + + assertEquals( + 7, + cache.getOrCreate("k") { + builds++ + 7 + }, + ) + assertEquals( + 7, + cache.getOrCreate("k") { + builds++ + 99 + }, + ) + + assertEquals(1, builds, "the builder must not run for a key that is already present") + assertEquals(7, cache.get("k")) + assertEquals(1, cache.size()) + } + + @Test + fun createIfAbsentReportsWhoInserted() { + val cache = LargeCache() + + assertTrue(cache.createIfAbsent("k") { 1 }, "the first call inserts") + assertFalse(cache.createIfAbsent("k") { 2 }, "the second call must not report an insert") + + assertEquals(1, cache.get("k"), "a losing createIfAbsent must not overwrite") + assertEquals(1, cache.size()) + } + + @Test + fun keysAndValuesSeeLaterWrites() { + val cache = cacheOf("a" to 1) + + // Reads the collections first, so an implementation that caches a snapshot has + // one to go stale. + assertEquals(setOf("a"), cache.keys().toSet()) + assertEquals(listOf(1), cache.values().toList()) + + cache.put("b", 2) + + assertEquals(setOf("a", "b"), cache.keys().toSet()) + assertEquals(listOf(1, 2), cache.values().sorted()) + + cache.remove("a") + + assertEquals(setOf("b"), cache.keys().toSet()) + assertEquals(listOf(2), cache.values().toList()) + } + + @Test + fun bulkReadsSeeLaterWrites() { + val cache = cacheOf("a" to 1) + + // Same idea for the collector path: warm every kind of bulk read, then mutate + // and re-read. Guards the lazily rebuilt snapshot in the linux actual. + assertEquals(1, cache.count { _, _ -> true }) + assertEquals(1, cache.sumOf { _, v -> v }) + + cache.put("b", 2) + assertEquals(2, cache.count { _, _ -> true }) + assertEquals(3, cache.sumOf { _, v -> v }) + + cache.put("a", 10) + assertEquals(12, cache.sumOf { _, v -> v }, "an overwrite must invalidate a cached read view") + + cache.remove("b") + assertEquals(10, cache.sumOf { _, v -> v }) + + cache.getOrCreate("c") { 5 } + assertEquals(15, cache.sumOf { _, v -> v }) + + cache.createIfAbsent("d") { 100 } + assertEquals(115, cache.sumOf { _, v -> v }) + + cache.clear() + assertEquals(0, cache.count { _, _ -> true }) + assertEquals(0, cache.sumOf { _, v -> v }) + } + + @Test + fun forEachVisitsEveryEntry() { + val cache = cacheOf("a" to 1, "b" to 2, "c" to 3) + + val seen = mutableMapOf() + cache.forEach { k, v -> seen[k] = v } + + assertEquals(mapOf("a" to 1, "b" to 2, "c" to 3), seen) + } + + @Test + fun filterAndCount() { + val cache = cacheOf("a" to 1, "b" to 2, "c" to 3, "d" to 4) + + assertEquals(listOf(2, 4), cache.filter { _, v -> v % 2 == 0 }.sorted()) + assertEquals(setOf(2, 4), cache.filterIntoSet { _, v -> v % 2 == 0 }) + assertEquals(2, cache.count { _, v -> v % 2 == 0 }) + assertEquals(1, cache.count { k, _ -> k == "a" }) + assertEquals(0, cache.count { _, _ -> false }) + } + + @Test + fun mapVariants() { + val cache = cacheOf("a" to 1, "b" to 2, "c" to 3) + + assertEquals(listOf(2, 4, 6), cache.map { _, v -> v * 2 }.sorted()) + assertEquals(listOf(2, 3), cache.mapNotNull { _, v -> if (v > 1) v else null }.sorted()) + assertEquals(setOf(2, 3), cache.mapNotNullIntoSet { _, v -> if (v > 1) v else null }) + assertEquals( + listOf(1, 1, 2, 2, 3, 3), + cache.mapFlatten { _, v -> listOf(v, v) }.sorted(), + ) + assertEquals( + setOf(1, 2, 3), + cache.mapFlattenIntoSet { _, v -> listOf(v, v) }, + ) + assertEquals( + listOf(2, 3), + cache.mapFlatten { _, v -> if (v > 1) listOf(v) else null }.sorted(), + "a null collection contributes nothing", + ) + } + + @Test + fun aggregates() { + val cache = cacheOf("a" to 1, "b" to 2, "c" to 3, "d" to 4) + + assertEquals(10, cache.sumOf { _, v -> v }) + assertEquals(10L, cache.sumOfLong { _, v -> v.toLong() }) + assertEquals(4, cache.maxOrNullOf({ _, _ -> true }, naturalOrder())) + assertEquals(3, cache.maxOrNullOf({ _, v -> v % 2 == 1 }, naturalOrder())) + assertNull(cache.maxOrNullOf({ _, _ -> false }, naturalOrder())) + } + + @Test + fun groupings() { + val cache = cacheOf("a" to 1, "b" to 2, "c" to 3, "d" to 4) + + assertEquals( + mapOf(0 to listOf(2, 4), 1 to listOf(1, 3)), + cache.groupBy { _, v -> v % 2 }.mapValues { it.value.sorted() }, + ) + assertEquals(mapOf(0 to 2, 1 to 2), cache.countByGroup { _, v -> v % 2 }) + assertEquals( + mapOf(0 to 6L, 1 to 4L), + cache.sumByGroup({ _, v -> v % 2 }, { _, v -> v.toLong() }), + ) + } + + @Test + fun associates() { + val cache = cacheOf("a" to 1, "b" to 2) + + assertEquals(mapOf(1 to "a", 2 to "b"), cache.associate { k, v -> v to k }) + assertEquals(mapOf("a" to 2, "b" to 4), cache.associateWith { _, v -> v * 2 }) + assertEquals(mapOf("a" to null, "b" to 4), cache.associateWith { _, v -> if (v > 1) v * 2 else null }) + } + + @Test + fun joinToStringRendersEveryEntry() { + val cache = cacheOf("a" to 1, "b" to 2, "c" to 3) + + val rendered = + cache.joinToString( + separator = ",", + prefix = "[", + postfix = "]", + limit = -1, + truncated = "...", + ) { k, v -> "$k=$v" } + + assertTrue(rendered.startsWith("[") && rendered.endsWith("]"), rendered) + assertEquals( + listOf("a=1", "b=2", "c=3"), + rendered + .removePrefix("[") + .removeSuffix("]") + .split(",") + .sorted(), + ) + + assertEquals("[]", LargeCache().joinToString(",", "[", "]", -1, "...") { k, v -> "$k=$v" }) + } + + @Test + fun insertingManyEntriesStaysLinear() { + // Regression guard for the copy-on-write linux actual this replaced, where every + // put rebuilt the whole map: this loop cost ~1.25 billion entry copies there and + // is instant against any O(1)-write implementation. No wall-clock assertion — + // the run time itself is the signal. + val n = 50_000 + val cache = LargeCache() + + for (i in 0 until n) { + cache.put(i, i) + // A full scan every thousandth insert, so an implementation that rebuilds a + // read view on write has to rebuild it ~50 times rather than amortize it away. + if (i % 1_000 == 0) assertEquals(i + 1, cache.count { _, _ -> true }) + } + + assertEquals(n, cache.size()) + assertEquals(0, cache.get(0)) + assertEquals(n - 1, cache.get(n - 1)) + assertEquals(n.toLong() * (n - 1) / 2, cache.sumOfLong { _, v -> v.toLong() }) + + for (i in 0 until n step 2) cache.remove(i) + + assertEquals(n / 2, cache.size()) + assertNull(cache.get(0)) + assertEquals(1, cache.get(1)) + } +} diff --git a/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/ConcurrentHashCache.linux.kt b/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/ConcurrentHashCache.linux.kt index 29d4d13487..1f697463eb 100644 --- a/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/ConcurrentHashCache.linux.kt +++ b/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/ConcurrentHashCache.linux.kt @@ -20,30 +20,37 @@ */ package com.vitorpamplona.quartz.utils.cache -import kotlin.concurrent.AtomicReference +import com.vitorpamplona.quartz.utils.concurrent.PlatformLock +import com.vitorpamplona.quartz.utils.concurrent.withLock -// Copy-on-write, mirroring LargeCache.linux: correct and simple; the linux -// target is CI-only so write cost is acceptable. +/** + * Linux/Native actual for [ConcurrentHashCache]. + * + * Was copy-on-write, mirroring the old `LargeCache.linux`: every [put] rebuilt the + * whole map under a CAS retry loop, so writes were O(n) and a decode burst against a + * warm cache was O(n^2). That is a bad shape for this class in particular — its only + * caller, `CachingEventDecoder`, writes once per event arriving from a relay. + * + * Now a plain [HashMap] guarded by a [PlatformLock]: O(1) writes, and no CAS retry + * to livelock under a write burst. Same lock choice and the same residual (a spin + * lock on this target) as `LargeCache.linux.kt` — see its docs. + */ actual class ConcurrentHashCache { - private val mapRef = AtomicReference(HashMap()) + private val lock = PlatformLock() + private val map = HashMap() - actual fun get(key: K): V? = mapRef.value[key] + actual fun get(key: K): V? = lock.withLock { map[key] } actual fun put( key: K, value: V, ) { - while (true) { - val current = mapRef.value - val copy = HashMap(current) - copy[key] = value - if (mapRef.compareAndSet(current, copy)) return - } + lock.withLock { map[key] = value } } - actual fun size(): Int = mapRef.value.size + actual fun size(): Int = lock.withLock { map.size } actual fun clear() { - mapRef.value = HashMap() + lock.withLock { map.clear() } } } diff --git a/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCache.linux.kt b/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCache.linux.kt index 2e49f22a5c..6887ef1bae 100644 --- a/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCache.linux.kt +++ b/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCache.linux.kt @@ -20,41 +20,112 @@ */ package com.vitorpamplona.quartz.utils.cache -import kotlin.concurrent.AtomicReference +import com.vitorpamplona.quartz.utils.concurrent.PlatformLock +import com.vitorpamplona.quartz.utils.concurrent.withLock +/** + * Linux/Native actual for [LargeCache] — the store behind Amethyst's `LocalCache`. + * + * ## Why this is not copy-on-write any more + * + * The first cut of this file kept a `LinkedHashMap` inside an `AtomicReference` and + * replaced it wholesale on every write. That made each [put] **O(n) in the size of + * the cache**: inserting n entries cost O(n^2) copies, so a cache holding 100k notes + * paid a 100k-entry map copy per arriving event. It was also *not* actually + * thread-safe — the read-copy-write was not a CAS loop, so two concurrent writers + * silently dropped one of the two writes. + * + * That went unnoticed because no CI job runs the linuxX64 target, so nothing ever + * pushed volume through this class. + * + * ## What it does instead + * + * A single mutable [LinkedHashMap] guarded by a [PlatformLock], plus a lazily built + * read snapshot: + * + * - **Point operations** ([get], [put], [remove], [containsKey], [size], …) take the + * lock, touch the live map, and return. O(1), no copying. + * - **Bulk operations** (`filter`/`map`/`forEach`/…) run against [cachedSnapshot], a + * point-in-time copy rebuilt on the first bulk call after a write and reused until + * the next write. So a copy costs O(n) at most once per write epoch, on operations + * that are already O(n), and a run of reads with no interleaved write copies + * nothing at all. + * + * The snapshot is what lets the bulk operations invoke caller-supplied lambdas + * **outside** the critical section. That matters: `PlatformLock` is not reentrant on + * this target, and `LocalCache` predicates routinely call back into the same cache — + * running them under the lock would self-deadlock. It also removes the + * `ConcurrentModificationException` window the previous `entries.toList()` dance was + * working around. + * + * ## Known residual + * + * Building the snapshot holds the lock for O(n), and the linux `PlatformLock` is a + * spin lock (Kotlin/Native ships no parking lock and there is no Foundation here — + * see `PlatformLock.linux.kt`). A writer racing a snapshot build therefore busy-waits + * for the duration of the copy. That is still strictly better than what it replaces — + * copy-on-write did the same O(n) copy on *every write* and lost concurrent ones — + * but it is the reason `PlatformLock.linux.kt` flags a pthread mutex as the next step + * if this target ever hosts a genuinely contended workload. + * + * Ordering note: iteration follows insertion order here, sorted-key order on + * JVM/Android (`ConcurrentSkipListMap`) and hash order on Apple. Nothing in the + * codebase depends on a specific order, and the `from`/`to` range overloads below + * degrade to a full scan on this target exactly as they did before — they have no + * callers outside the JVM-only `LargeSoftCache`. + */ actual class LargeCache : ICacheOperations { - private val mapRef = AtomicReference(LinkedHashMap()) + private val lock = PlatformLock() - private inline fun withMap(block: (LinkedHashMap) -> R): R = block(mapRef.value) + /** The live store. Every access must hold [lock]. */ + private val map = LinkedHashMap() - private inline fun mutate(block: (LinkedHashMap) -> Unit) { - val copy = LinkedHashMap(mapRef.value) - block(copy) - mapRef.value = copy - } + /** + * A copy of [map] handed to bulk operations so they can run caller lambdas + * without holding [lock]. Null means "stale, rebuild on next use". Guarded by + * [lock]; never mutated once published, so readers may keep it as long as they + * like. + */ + private var cachedSnapshot: Map? = null - actual fun keys(): Set = withMap { LinkedHashSet(it.keys) } - - actual fun values(): Iterable = withMap { ArrayList(it.values) } - - actual fun get(key: K): V? = withMap { it[key] } - - actual fun remove(key: K): V? { - val current = mapRef.value - val value = current[key] - if (value != null) { - mutate { it.remove(key) } + private fun snapshot(): Map = + lock.withLock { + cachedSnapshot ?: LinkedHashMap(map).also { cachedSnapshot = it } } - return value - } - actual fun isEmpty(): Boolean = withMap { it.isEmpty() } + /** Runs [block] over a stable snapshot, outside the lock. */ + private inline fun withMap(block: (Map) -> R): R = block(snapshot()) + + /** Runs [block] over the live map under [lock] without invalidating the snapshot. */ + private inline fun read(block: (MutableMap) -> R): R = lock.withLock { block(map) } + + /** Runs [block] over the live map under [lock] and drops the read snapshot. */ + private inline fun mutate(block: (MutableMap) -> R): R = + lock.withLock { + cachedSnapshot = null + block(map) + } + + actual fun keys(): Set = snapshot().keys + + actual fun values(): Iterable = snapshot().values + + actual fun get(key: K): V? = read { it[key] } + + actual fun remove(key: K): V? = + lock.withLock { + val removed = map.remove(key) + if (removed != null) cachedSnapshot = null + removed + } + + actual fun isEmpty(): Boolean = read { it.isEmpty() } actual fun clear() { - mapRef.value = LinkedHashMap() + mutate { it.clear() } } - actual fun containsKey(key: K): Boolean = withMap { it.containsKey(key) } + actual fun containsKey(key: K): Boolean = read { it.containsKey(key) } actual fun put( key: K, @@ -63,33 +134,61 @@ actual class LargeCache : ICacheOperations { mutate { it[key] = value } } + /** + * Mirrors the JVM actual's `putIfAbsent`: [builder] runs outside the lock (it is + * caller code and must not be able to re-enter a non-reentrant lock), and the + * insert is only published if no one won the race in the meantime. + */ actual fun getOrCreate( key: K, builder: (key: K) -> V, ): V { - val existing = get(key) - if (existing != null) return existing + read { it[key] }?.let { return it } + val newObject = builder(key) - mutate { it[key] = newObject } - return get(key) ?: newObject + + return lock.withLock { + val existing = map[key] + if (existing != null) { + existing + } else { + map[key] = newObject + cachedSnapshot = null + newObject + } + } } + /** + * True only when *this* call inserted the value — matching the JVM actual's + * `putIfAbsent(key, newObject) == null`. The previous implementation returned + * `get(key) != null`, which also reported true when another thread had just + * created the entry, double-firing whatever the caller does with a fresh key. + */ actual fun createIfAbsent( key: K, builder: (key: K) -> V, ): Boolean { - val existing = get(key) - if (existing != null) return false + if (read { it.containsKey(key) }) return false + val newObject = builder(key) - mutate { it[key] = newObject } - return get(key) != null + + return lock.withLock { + if (map.containsKey(key)) { + false + } else { + map[key] = newObject + cachedSnapshot = null + true + } + } } - actual override fun size(): Int = withMap { it.size } + actual override fun size(): Int = read { it.size } actual override fun forEach(consumer: ICacheBiConsumer) { - // Snapshot entries to avoid ConcurrentModificationException - withMap { map -> map.entries.toList() }.forEach { consumer.accept(it.key, it.value) } + // The snapshot is already immutable, so no defensive entries.toList() is needed. + withMap { map -> map.forEach { consumer.accept(it.key, it.value) } } } actual override fun filter(consumer: CacheCollectors.BiFilter): List = withMap { map -> map.filter { consumer.filter(it.key, it.value) }.values.toList() } diff --git a/quartz/src/linuxTest/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCacheConcurrencyTest.kt b/quartz/src/linuxTest/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCacheConcurrencyTest.kt new file mode 100644 index 0000000000..07b674700c --- /dev/null +++ b/quartz/src/linuxTest/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCacheConcurrencyTest.kt @@ -0,0 +1,120 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.quartz.utils.cache + +import kotlin.native.concurrent.ObsoleteWorkersApi +import kotlin.native.concurrent.TransferMode +import kotlin.native.concurrent.Worker +import kotlin.test.Test +import kotlin.test.assertEquals + +/** + * The linux actual is the only [LargeCache] whose thread safety is hand-rolled rather + * than delegated to a concurrent map, so it gets its own multi-threaded test. The + * copy-on-write version this replaced would fail every assertion here: its + * read-copy-write was not a CAS loop, so concurrent writers silently dropped each + * other's entries. + * + * Uses `Worker` rather than coroutines on purpose — a coroutine dispatcher gives no + * guarantee of genuine parallelism, and parallelism is the whole point. + */ +@OptIn(ObsoleteWorkersApi::class) +class LargeCacheConcurrencyTest { + private val workerCount = 4 + private val perWorker = 5_000 + + private fun inParallel(job: (workerId: Int) -> R): List { + val workers = List(workerCount) { Worker.start() } + val futures = + workers.mapIndexed { id, worker -> + worker.execute(TransferMode.SAFE, { Pair(job, id) }) { (block, workerId) -> block(workerId) } + } + val results = futures.map { it.result } + workers.forEach { it.requestTermination().result } + return results + } + + @Test + fun concurrentPutsKeepEveryEntry() { + val cache = LargeCache() + + inParallel { id -> + repeat(perWorker) { i -> cache.put(id * perWorker + i, i) } + } + + assertEquals(workerCount * perWorker, cache.size()) + for (id in 0 until workerCount) { + assertEquals(0, cache.get(id * perWorker)) + assertEquals(perWorker - 1, cache.get(id * perWorker + perWorker - 1)) + } + } + + @Test + fun concurrentGetOrCreateBuildsOneValuePerKey() { + val cache = LargeCache() + + // Every worker races for the same 500 keys. Whoever wins, all of them must end + // up holding the identical instance, and the cache must hold exactly 500. + val seen = + inParallel { _ -> + (0 until 500).map { key -> cache.getOrCreate(key) { "v$it" } } + } + + assertEquals(500, cache.size()) + seen.forEach { perWorkerValues -> + perWorkerValues.forEachIndexed { key, value -> + assertEquals(cache.get(key), value, "getOrCreate handed out a value it did not publish") + } + } + } + + @Test + fun concurrentCreateIfAbsentReportsExactlyOneInsertPerKey() { + val cache = LargeCache() + + val insertsPerWorker = inParallel { _ -> (0 until 500).count { key -> cache.createIfAbsent(key) { key } } } + + assertEquals(500, cache.size()) + assertEquals(500, insertsPerWorker.sum(), "exactly one caller per key may report an insert") + } + + @Test + fun bulkReadsStayConsistentWhileWritesLand() { + val cache = LargeCache() + repeat(1_000) { cache.put(it, it) } + + // Writers churn the map while readers walk snapshots of it. A reader must never + // see a torn map, and must never crash on a concurrent modification. + inParallel { id -> + if (id % 2 == 0) { + repeat(perWorker) { i -> cache.put(1_000 + id * perWorker + i, i) } + } else { + repeat(200) { + val sum = cache.sumOfLong { _, v -> v.toLong() } + check(sum >= 0) { "unexpected negative sum $sum" } + cache.count { _, _ -> true } + } + } + } + + assertEquals(1_000 + (workerCount / 2) * perWorker, cache.size()) + } +} diff --git a/quartz/src/linuxTest/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCacheRangeFallbackTest.kt b/quartz/src/linuxTest/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCacheRangeFallbackTest.kt new file mode 100644 index 0000000000..1ac079bb92 --- /dev/null +++ b/quartz/src/linuxTest/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCacheRangeFallbackTest.kt @@ -0,0 +1,73 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.quartz.utils.cache + +import kotlin.test.Test +import kotlin.test.assertContentEquals +import kotlin.test.assertEquals + +/** + * The `from`/`to` overloads narrow to a key range only where the backing map is sorted + * (`ConcurrentSkipListMap` on JVM/Android). This target's map is not sorted, and `K` + * carries no `Comparable` bound to sort it by, so every range overload deliberately + * degrades to a full scan and leaves the caller's predicate to do the filtering. + * + * That is safe today because the range overloads have no callers outside the JVM-only + * `LargeSoftCache`. This test pins the behaviour so the fallback stays *consistent* + * with the unbounded form — the property anything sharing code across targets would + * rely on — rather than silently returning a different set on this platform. + * + * Linux-only on purpose: it is a statement about this actual. The Apple actual walks + * its map by index instead, which is a different (and order-dependent) contract. + */ +class LargeCacheRangeFallbackTest { + private val cache = + LargeCache().apply { + put("a", 1) + put("b", 2) + put("c", 3) + } + + private val all = CacheCollectors.BiFilter { _, _ -> true } + + @Test + fun rangeOverloadsMatchTheirUnboundedForm() { + assertContentEquals(cache.filter(all), cache.filter("a", "c", all)) + assertEquals(cache.filterIntoSet(all), cache.filterIntoSet("a", "c", all)) + assertEquals(cache.count(all), cache.count("a", "c", all)) + assertEquals(cache.sumOf { _, v -> v }, cache.sumOf("a", "c") { _, v -> v }) + assertEquals(cache.sumOfLong { _, v -> v.toLong() }, cache.sumOfLong("a", "c") { _, v -> v.toLong() }) + assertContentEquals(cache.map { _, v -> v }, cache.map("a", "c") { _, v -> v }) + assertContentEquals(cache.mapNotNull { _, v -> v }, cache.mapNotNull("a", "c") { _, v -> v }) + assertEquals(cache.associate { k, v -> k to v }, cache.associate("a", "c") { k, v -> k to v }) + assertEquals(cache.associateWith { _, v -> v }, cache.associateWith("a", "c") { _, v -> v }) + assertEquals(cache.countByGroup { _, v -> v % 2 }, cache.countByGroup("a", "c") { _, v -> v % 2 }) + } + + @Test + fun aNarrowerRangeStillScansEverything() { + // Documents the degradation explicitly: on a sorted target this would return + // only "a", here it returns all three. Anyone who later gives this actual real + // range support should update this test rather than discover the difference in + // production. + assertEquals(3, cache.count("a", "a", all)) + } +} diff --git a/quartz/src/nativeMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.native.kt b/quartz/src/nativeMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.native.kt index 8805e8d5a7..4e81970190 100644 --- a/quartz/src/nativeMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.native.kt +++ b/quartz/src/nativeMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.native.kt @@ -23,10 +23,12 @@ package com.vitorpamplona.quartz.utils.concurrent import kotlin.concurrent.atomics.AtomicReference import kotlin.concurrent.atomics.ExperimentalAtomicApi -// Copy-on-write, mirroring ConcurrentHashCache.linux: correct and simple. The -// native targets never run the crawl this backs (it is JVM/Android-only work); -// they only compile it, so the O(n)-per-write cost is irrelevant. A CAS retry -// loop gives getOrPut/merge the same atomicity the JVM actual gets for free. +// Copy-on-write: correct and simple. Unlike LargeCache/ConcurrentHashCache — which +// back LocalCache and the event decoder and were moved off copy-on-write for exactly +// this reason — the native targets never run the crawl this backs (it is +// JVM/Android-only work); they only compile it, so the O(n)-per-write cost is +// irrelevant. A CAS retry loop gives getOrPut/merge the same atomicity the JVM actual +// gets for free. @OptIn(ExperimentalAtomicApi::class) actual class ConcurrentMap { private val ref = AtomicReference(HashMap()) From 65ea34342e06ece3d33468a9cba0bb34729c5f9b Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 1 Sep 2026 17:57:22 +0000 Subject: [PATCH 2/6] perf(quartz): make linuxX64's LargeCache lock-free, not lock-based MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Follow-up to the previous commit, which fixed the O(n) write by putting a PlatformLock around a mutable map. That traded one problem for another: the JVM/Android actual is a ConcurrentSkipListMap, where readers never block and writers publish with a CAS, and a global lock is a step down from that — worse, the linux PlatformLock is a spin lock, so a reader could burn a core waiting on a writer that had been descheduled. Keep copy-on-write's shape instead — an immutable map behind an AtomicReference, which is what made reads free in the first place — and fix the two things that were actually wrong with it. Copying a LinkedHashMap is O(n); a HAMT's putting() shares structure and copies only the path to the changed key, O(log32 n). And the read-copy-write was not a CAS loop, so concurrent writers dropped each other's entries; now they retry. Reads (get/containsKey/size/keys/values) are a single atomic load plus a lookup. Bulk operations iterate that same immutable map with no copy, so caller lambdas run outside any critical section and a LocalCache predicate that reaches back into the cache cannot deadlock. Writes are a CAS retry. This is already the house pattern for shared mutable state in commonMain — FilterIndex and nip86 BanStore hold state in one AtomicReference over persistent collections and mutate it with the same loop — and kotlinx-collections-immutable is already a quartz commonMain dependency. Measured on linuxX64 (-opt, ms per loop), vs copy-on-write and vs the lock variant this replaces: n=20,000 fill reads 20 scans mixed copy-on-write 17,949 2 13 25,278 lock+HashMap 1 0 12 35 HAMT+CAS 14 0 19 28 n=200,000 fill reads 20 scans mixed lock+HashMap 44 9 177 6,736 HAMT+CAS 197 12 237 2,486 Write-only, the lock wins ~4x. But LocalCache interleaves full-cache scans with arriving events, and there the lock must rebuild an O(n) read snapshot per write epoch: it loses by 2.7x at 200k. So the non-blocking design also wins the workload that matters. ConcurrentHashCache.linux gets the same treatment; iteration order becomes hash order (as on Apple) rather than insertion order. Nothing depends on it. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01HxQ1QuyzSkR38iFHbREjoS --- .../utils/cache/ConcurrentHashCache.linux.kt | 37 +-- .../quartz/utils/cache/LargeCache.linux.kt | 230 +++++++++--------- 2 files changed, 141 insertions(+), 126 deletions(-) diff --git a/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/ConcurrentHashCache.linux.kt b/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/ConcurrentHashCache.linux.kt index 1f697463eb..d59653313b 100644 --- a/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/ConcurrentHashCache.linux.kt +++ b/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/ConcurrentHashCache.linux.kt @@ -20,37 +20,44 @@ */ package com.vitorpamplona.quartz.utils.cache -import com.vitorpamplona.quartz.utils.concurrent.PlatformLock -import com.vitorpamplona.quartz.utils.concurrent.withLock +import kotlinx.collections.immutable.PersistentMap +import kotlinx.collections.immutable.persistentHashMapOf +import kotlin.concurrent.atomics.AtomicReference +import kotlin.concurrent.atomics.ExperimentalAtomicApi /** * Linux/Native actual for [ConcurrentHashCache]. * - * Was copy-on-write, mirroring the old `LargeCache.linux`: every [put] rebuilt the - * whole map under a CAS retry loop, so writes were O(n) and a decode burst against a - * warm cache was O(n^2). That is a bad shape for this class in particular — its only - * caller, `CachingEventDecoder`, writes once per event arriving from a relay. + * Same fix, and for the same reason, as `LargeCache.linux.kt` — read its docs for the + * full rationale. This was copy-on-write over a plain `HashMap`, so every [put] rebuilt + * the entire map: O(n) per write, and a CAS retry loop that re-did the whole rebuild on + * every lost race. Bad anywhere; worst here, because the only caller is + * `CachingEventDecoder`, which writes once per event arriving from a relay. * - * Now a plain [HashMap] guarded by a [PlatformLock]: O(1) writes, and no CAS retry - * to livelock under a write burst. Same lock choice and the same residual (a spin - * lock on this target) as `LargeCache.linux.kt` — see its docs. + * A HAMT keeps the wait-free single-load read and the non-blocking write while making + * the write O(log32 n) — [PersistentMap.putting] shares structure with the map it came + * from and copies only the path to the changed key. */ +@OptIn(ExperimentalAtomicApi::class) actual class ConcurrentHashCache { - private val lock = PlatformLock() - private val map = HashMap() + private val ref = AtomicReference>(persistentHashMapOf()) - actual fun get(key: K): V? = lock.withLock { map[key] } + actual fun get(key: K): V? = ref.load()[key] actual fun put( key: K, value: V, ) { - lock.withLock { map[key] = value } + while (true) { + val current = ref.load() + val next = current.putting(key, value) + if (next === current || ref.compareAndSet(current, next)) return + } } - actual fun size(): Int = lock.withLock { map.size } + actual fun size(): Int = ref.load().size actual fun clear() { - lock.withLock { map.clear() } + ref.store(persistentHashMapOf()) } } diff --git a/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCache.linux.kt b/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCache.linux.kt index 6887ef1bae..51251fb309 100644 --- a/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCache.linux.kt +++ b/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCache.linux.kt @@ -20,174 +20,182 @@ */ package com.vitorpamplona.quartz.utils.cache -import com.vitorpamplona.quartz.utils.concurrent.PlatformLock -import com.vitorpamplona.quartz.utils.concurrent.withLock +import kotlinx.collections.immutable.PersistentMap +import kotlinx.collections.immutable.persistentHashMapOf +import kotlin.concurrent.atomics.AtomicReference +import kotlin.concurrent.atomics.ExperimentalAtomicApi /** * Linux/Native actual for [LargeCache] — the store behind Amethyst's `LocalCache`. * - * ## Why this is not copy-on-write any more + * Mirrors the progress guarantees of the JVM/Android actual (`ConcurrentSkipListMap`): + * **readers never block and never wait for a writer**, and writers publish with a CAS + * rather than by holding a lock. Nothing here can be descheduled while excluding + * everyone else, which is the failure mode `PlatformLock`'s docs describe. * - * The first cut of this file kept a `LinkedHashMap` inside an `AtomicReference` and - * replaced it wholesale on every write. That made each [put] **O(n) in the size of - * the cache**: inserting n entries cost O(n^2) copies, so a cache holding 100k notes - * paid a 100k-entry map copy per arriving event. It was also *not* actually - * thread-safe — the read-copy-write was not a CAS loop, so two concurrent writers - * silently dropped one of the two writes. + * ## Why it is shaped this way * - * That went unnoticed because no CI job runs the linuxX64 target, so nothing ever - * pushed volume through this class. + * The first cut kept a `LinkedHashMap` inside an `AtomicReference` and replaced it + * wholesale on every write. The immutable-snapshot *idea* was right — it is what makes + * reads free — but two things were wrong with it: copying a `LinkedHashMap` makes each + * [put] **O(n) in the size of the cache** (filling n entries costs O(n^2), so a cache + * holding 100k notes paid a 100k-entry copy per arriving event), and the + * read-copy-write was not a CAS loop, so two concurrent writers silently dropped one + * of the two writes. * - * ## What it does instead + * Both fall away by swapping the map for a HAMT. [PersistentMap.putting] shares + * structure with the map it came from and only copies the nodes on the path to the + * changed key — O(log32 n), so ~4 small array copies at a million entries instead of a + * million-entry rehash — and the CAS loop makes concurrent writers retry instead of + * clobbering each other. * - * A single mutable [LinkedHashMap] guarded by a [PlatformLock], plus a lazily built - * read snapshot: + * So: + * - **Reads** ([get], [containsKey], [size], [keys], [values]) are a single atomic load + * plus a lookup on an immutable map. No lock, no allocation, no copy. + * - **Bulk operations** (`filter`/`map`/`forEach`/…) iterate that same immutable map + * directly. No defensive `entries.toList()`, no `ConcurrentModificationException` + * window, and — because no lock is held while a caller's lambda runs — no way for a + * `LocalCache` predicate that reaches back into the cache to deadlock. + * - **Writes** are a CAS retry loop over a structurally shared copy. * - * - **Point operations** ([get], [put], [remove], [containsKey], [size], …) take the - * lock, touch the live map, and return. O(1), no copying. - * - **Bulk operations** (`filter`/`map`/`forEach`/…) run against [cachedSnapshot], a - * point-in-time copy rebuilt on the first bulk call after a write and reused until - * the next write. So a copy costs O(n) at most once per write epoch, on operations - * that are already O(n), and a run of reads with no interleaved write copies - * nothing at all. + * ## What it costs * - * The snapshot is what lets the bulk operations invoke caller-supplied lambdas - * **outside** the critical section. That matters: `PlatformLock` is not reentrant on - * this target, and `LocalCache` predicates routinely call back into the same cache — - * running them under the lock would self-deadlock. It also removes the - * `ConcurrentModificationException` window the previous `entries.toList()` dance was - * working around. + * A write costs more than a `HashMap.put` under a lock would — a few node copies rather + * than one bucket store. Measured on this target (linuxX64, `-opt`, ms for the whole + * loop) against the copy-on-write version this replaces and against a + * `PlatformLock` + `HashMap` variant that was the other candidate: * - * ## Known residual + * ``` + * n=20,000 fill reads 20 scans mixed (put + scan every 1k) + * copy-on-write 17,949 2 13 25,278 + * lock + HashMap 1 0 12 35 + * HAMT + CAS 14 0 19 28 * - * Building the snapshot holds the lock for O(n), and the linux `PlatformLock` is a - * spin lock (Kotlin/Native ships no parking lock and there is no Foundation here — - * see `PlatformLock.linux.kt`). A writer racing a snapshot build therefore busy-waits - * for the duration of the copy. That is still strictly better than what it replaces — - * copy-on-write did the same O(n) copy on *every write* and lost concurrent ones — - * but it is the reason `PlatformLock.linux.kt` flags a pthread mutex as the next step - * if this target ever hosts a genuinely contended workload. + * n=200,000 fill reads 20 scans mixed + * lock + HashMap 44 9 177 6,736 + * HAMT + CAS 197 12 237 2,486 + * ``` * - * Ordering note: iteration follows insertion order here, sorted-key order on - * JVM/Android (`ConcurrentSkipListMap`) and hash order on Apple. Nothing in the - * codebase depends on a specific order, and the `from`/`to` range overloads below - * degrade to a full scan on this target exactly as they did before — they have no - * callers outside the JVM-only `LargeSoftCache`. + * Writing nothing but writes, the lock wins ~4x. But that is not the shape `LocalCache` + * has: it interleaves scans (every feed filters the whole cache) with arriving events, + * and there the lock has to rebuild an O(n) read snapshot after each write epoch, so it + * loses by ~2.7x at 200k. Reads and scans are close either way. The lock-free version + * therefore wins the workload that matters *and* is the one that never blocks a reader. + * + * This is the house pattern for shared mutable state in `commonMain` already — see + * `nip01Core.relay.filters.FilterIndex` and `nip86RelayManagement.server.BanStore`, + * both of which hold their state in one `AtomicReference` over persistent collections + * and mutate it with the same CAS loop. + * + * ## Notes + * + * Anything that reads twice — [remove], [getOrCreate], [createIfAbsent] — retries on a + * lost CAS rather than locking, so the pair is atomic without excluding readers. + * [clear] publishes an empty map unconditionally and can therefore drop a write that + * lands concurrently, exactly as `ConcurrentSkipListMap.clear()` can. + * + * Iteration order is hash order, as on Apple; JVM/Android is sorted-key order + * (`ConcurrentSkipListMap`). Nothing in the codebase depends on a specific order. The + * `from`/`to` range overloads below degrade to a full scan on this target, as they + * always have — they have no callers outside the JVM-only `LargeSoftCache`. */ +@OptIn(ExperimentalAtomicApi::class) actual class LargeCache : ICacheOperations { - private val lock = PlatformLock() - - /** The live store. Every access must hold [lock]. */ - private val map = LinkedHashMap() + private val ref = AtomicReference>(persistentHashMapOf()) /** - * A copy of [map] handed to bulk operations so they can run caller lambdas - * without holding [lock]. Null means "stale, rebuild on next use". Guarded by - * [lock]; never mutated once published, so readers may keep it as long as they - * like. + * Runs [block] against an immutable point-in-time map. Outside any critical + * section, so [block] may call back into this cache freely. */ - private var cachedSnapshot: Map? = null + private inline fun withMap(block: (Map) -> R): R = block(ref.load()) - private fun snapshot(): Map = - lock.withLock { - cachedSnapshot ?: LinkedHashMap(map).also { cachedSnapshot = it } + /** + * Publishes [transform] of the current map with a CAS, retrying if a concurrent + * writer won. [transform] must be pure — it can run more than once — and returning + * the map it was given means "no change", which skips the CAS entirely. + */ + private inline fun mutate(transform: (PersistentMap) -> PersistentMap) { + while (true) { + val current = ref.load() + val next = transform(current) + if (next === current || ref.compareAndSet(current, next)) return } - - /** Runs [block] over a stable snapshot, outside the lock. */ - private inline fun withMap(block: (Map) -> R): R = block(snapshot()) - - /** Runs [block] over the live map under [lock] without invalidating the snapshot. */ - private inline fun read(block: (MutableMap) -> R): R = lock.withLock { block(map) } - - /** Runs [block] over the live map under [lock] and drops the read snapshot. */ - private inline fun mutate(block: (MutableMap) -> R): R = - lock.withLock { - cachedSnapshot = null - block(map) - } - - actual fun keys(): Set = snapshot().keys - - actual fun values(): Iterable = snapshot().values - - actual fun get(key: K): V? = read { it[key] } - - actual fun remove(key: K): V? = - lock.withLock { - val removed = map.remove(key) - if (removed != null) cachedSnapshot = null - removed - } - - actual fun isEmpty(): Boolean = read { it.isEmpty() } - - actual fun clear() { - mutate { it.clear() } } - actual fun containsKey(key: K): Boolean = read { it.containsKey(key) } + actual fun keys(): Set = ref.load().keys + + actual fun values(): Iterable = ref.load().values + + actual fun get(key: K): V? = ref.load()[key] + + actual fun remove(key: K): V? { + while (true) { + val current = ref.load() + val previous = current[key] ?: return null + if (ref.compareAndSet(current, current.removing(key))) return previous + } + } + + actual fun isEmpty(): Boolean = ref.load().isEmpty() + + actual fun clear() { + ref.store(persistentHashMapOf()) + } + + actual fun containsKey(key: K): Boolean = ref.load().containsKey(key) actual fun put( key: K, value: V, ) { - mutate { it[key] = value } + mutate { it.putting(key, value) } } /** - * Mirrors the JVM actual's `putIfAbsent`: [builder] runs outside the lock (it is - * caller code and must not be able to re-enter a non-reentrant lock), and the - * insert is only published if no one won the race in the meantime. + * Mirrors the JVM actual's `putIfAbsent`: [builder] runs at most once — outside the + * retry loop, since it is caller code — and the value is only published if no one + * won the race in the meantime. */ actual fun getOrCreate( key: K, builder: (key: K) -> V, ): V { - read { it[key] }?.let { return it } + ref.load()[key]?.let { return it } val newObject = builder(key) - return lock.withLock { - val existing = map[key] - if (existing != null) { - existing - } else { - map[key] = newObject - cachedSnapshot = null - newObject - } + while (true) { + val current = ref.load() + current[key]?.let { return it } + if (ref.compareAndSet(current, current.putting(key, newObject))) return newObject } } /** * True only when *this* call inserted the value — matching the JVM actual's * `putIfAbsent(key, newObject) == null`. The previous implementation returned - * `get(key) != null`, which also reported true when another thread had just - * created the entry, double-firing whatever the caller does with a fresh key. + * `get(key) != null`, which also reported true when another thread had just created + * the entry, double-firing whatever the caller does with a fresh key. */ actual fun createIfAbsent( key: K, builder: (key: K) -> V, ): Boolean { - if (read { it.containsKey(key) }) return false + if (ref.load().containsKey(key)) return false val newObject = builder(key) - return lock.withLock { - if (map.containsKey(key)) { - false - } else { - map[key] = newObject - cachedSnapshot = null - true - } + while (true) { + val current = ref.load() + if (current.containsKey(key)) return false + if (ref.compareAndSet(current, current.putting(key, newObject))) return true } } - actual override fun size(): Int = read { it.size } + actual override fun size(): Int = ref.load().size actual override fun forEach(consumer: ICacheBiConsumer) { - // The snapshot is already immutable, so no defensive entries.toList() is needed. + // The map is immutable, so this iterates a stable snapshot with no copy. withMap { map -> map.forEach { consumer.accept(it.key, it.value) } } } From 5e662fef0b908f6e33f9d837328588766de4f106 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 1 Sep 2026 18:36:33 +0000 Subject: [PATCH 3/6] perf(quartz): back linuxX64's LargeCache with a striped hash table MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The HAMT was the wrong structure for this workload. LocalCache fills on the order of 100,000 entries in a few seconds, and a persistent map allocates a fresh path of ~4-5 nodes for every write — including overwrites, which change no structure at all — then discards it. Over a 100k fill plus scans that is 24 GC cycles. Replace it with a chained hash table: lock-free reads, striped-lock writes. This is ConcurrentHashMap's shape, which Kotlin/Native does not ship. Adding a key prepends one node; overwriting one is a single volatile store into the node already there, allocating nothing; a scan walks the buckets in place. Chain nodes hold `next` immutably so a reader never sees it change, which is what lets reads take no lock at all — structural edits publish a new bucket head, and a resize rebuilds nodes rather than relinking them. Measured on linuxX64 (-opt), 100,000 String keys of event-id length: fill overwrite reads 20 scans mixed GCs heap HAMT + CAS 70 78 6 117 676 24 67MB lock + HashMap 13 6 3 71 1197 36 51MB striped 15 3 3 13 64 1 43MB "mixed" is a full fill with a whole-table scan every 1000 writes — the shape LocalCache actually has. Copy-on-write, the original, is off the scale: 20k entries alone took 18s to fill. Every bulk operation now walks the table directly instead of a snapshot, so scans allocate nothing beyond the result and caller lambdas run outside any critical section — a LocalCache predicate that reaches back into the cache cannot deadlock, and there is no ConcurrentModificationException window. getOrCreate and createIfAbsent are now the JVM actual's bodies verbatim over the same putIfAbsent contract. Honest difference from the JVM actual: ConcurrentSkipListMap is fully non-blocking, whereas writers here block writers hashing to the same one of 16 stripes, for a bucket walk of a few nodes. ConcurrentHashMap makes the same trade. Readers block for nothing. Adds LargeCacheCollisionTest, which forces every key into one bucket so the chain paths — in particular removal, which clones the nodes ahead of the target onto its tail — run deterministically rather than only on a chance collision. ConcurrentHashCache.linux moves onto the same table. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01HxQ1QuyzSkR38iFHbREjoS --- .../utils/cache/ConcurrentHashCache.linux.kt | 36 +- .../quartz/utils/cache/LargeCache.linux.kt | 346 ++++++++---------- .../quartz/utils/cache/StripedHashMap.kt | 309 ++++++++++++++++ .../utils/cache/LargeCacheCollisionTest.kt | 140 +++++++ 4 files changed, 605 insertions(+), 226 deletions(-) create mode 100644 quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/StripedHashMap.kt create mode 100644 quartz/src/linuxTest/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCacheCollisionTest.kt diff --git a/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/ConcurrentHashCache.linux.kt b/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/ConcurrentHashCache.linux.kt index d59653313b..294c96f168 100644 --- a/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/ConcurrentHashCache.linux.kt +++ b/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/ConcurrentHashCache.linux.kt @@ -20,44 +20,30 @@ */ package com.vitorpamplona.quartz.utils.cache -import kotlinx.collections.immutable.PersistentMap -import kotlinx.collections.immutable.persistentHashMapOf -import kotlin.concurrent.atomics.AtomicReference -import kotlin.concurrent.atomics.ExperimentalAtomicApi - /** - * Linux/Native actual for [ConcurrentHashCache]. + * Linux/Native actual for [ConcurrentHashCache], over the same [StripedHashMap] as + * `LargeCache.linux.kt` — read its docs for why. * - * Same fix, and for the same reason, as `LargeCache.linux.kt` — read its docs for the - * full rationale. This was copy-on-write over a plain `HashMap`, so every [put] rebuilt - * the entire map: O(n) per write, and a CAS retry loop that re-did the whole rebuild on - * every lost race. Bad anywhere; worst here, because the only caller is - * `CachingEventDecoder`, which writes once per event arriving from a relay. - * - * A HAMT keeps the wait-free single-load read and the non-blocking write while making - * the write O(log32 n) — [PersistentMap.putting] shares structure with the map it came - * from and copies only the path to the changed key. + * This one was the worst-placed of the copy-on-write caches: its only caller, + * `CachingEventDecoder`, writes once per event arriving from a relay, so every decode + * rebuilt the whole map under a CAS retry loop. Now a write touches one bucket and a + * read takes no lock at all. */ -@OptIn(ExperimentalAtomicApi::class) actual class ConcurrentHashCache { - private val ref = AtomicReference>(persistentHashMapOf()) + private val cache = StripedHashMap() - actual fun get(key: K): V? = ref.load()[key] + actual fun get(key: K): V? = cache.get(key) actual fun put( key: K, value: V, ) { - while (true) { - val current = ref.load() - val next = current.putting(key, value) - if (next === current || ref.compareAndSet(current, next)) return - } + cache.put(key, value) } - actual fun size(): Int = ref.load().size + actual fun size(): Int = cache.size() actual fun clear() { - ref.store(persistentHashMapOf()) + cache.clear() } } diff --git a/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCache.linux.kt b/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCache.linux.kt index 51251fb309..a718aee614 100644 --- a/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCache.linux.kt +++ b/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCache.linux.kt @@ -20,160 +20,86 @@ */ package com.vitorpamplona.quartz.utils.cache -import kotlinx.collections.immutable.PersistentMap -import kotlinx.collections.immutable.persistentHashMapOf -import kotlin.concurrent.atomics.AtomicReference -import kotlin.concurrent.atomics.ExperimentalAtomicApi - /** * Linux/Native actual for [LargeCache] — the store behind Amethyst's `LocalCache`. * - * Mirrors the progress guarantees of the JVM/Android actual (`ConcurrentSkipListMap`): - * **readers never block and never wait for a writer**, and writers publish with a CAS - * rather than by holding a lock. Nothing here can be descheduled while excluding - * everyone else, which is the failure mode `PlatformLock`'s docs describe. + * All of the concurrency and the performance rationale lives in [StripedHashMap]; this + * class is only the [ICacheOperations] surface over it. The short version: `LocalCache` + * fills ~100,000 entries in a few seconds while feeds scan the whole cache, so the + * backing store has to take writes at O(1) with next to no garbage *and* let a scan run + * in place without copying. A chained hash table with lock-free reads and striped-lock + * writes — `ConcurrentHashMap`'s shape, which Kotlin/Native does not ship — is the + * structure that does both; copy-on-write and a persistent HAMT each fail one half. * - * ## Why it is shaped this way + * Every bulk operation below walks the table through [StripedHashMap.forEachEntry], + * which takes no lock and allocates nothing beyond the result being built. Two things + * follow, both of which the earlier implementations had to work around: * - * The first cut kept a `LinkedHashMap` inside an `AtomicReference` and replaced it - * wholesale on every write. The immutable-snapshot *idea* was right — it is what makes - * reads free — but two things were wrong with it: copying a `LinkedHashMap` makes each - * [put] **O(n) in the size of the cache** (filling n entries costs O(n^2), so a cache - * holding 100k notes paid a 100k-entry copy per arriving event), and the - * read-copy-write was not a CAS loop, so two concurrent writers silently dropped one - * of the two writes. + * - The caller's lambda never runs inside a critical section, so a `LocalCache` + * predicate that reaches back into the cache cannot deadlock. + * - There is no snapshot and no defensive `entries.toList()`, so no + * `ConcurrentModificationException` window and no per-scan copy. * - * Both fall away by swapping the map for a HAMT. [PersistentMap.putting] shares - * structure with the map it came from and only copies the nodes on the path to the - * changed key — O(log32 n), so ~4 small array copies at a million entries instead of a - * million-entry rehash — and the CAS loop makes concurrent writers retry instead of - * clobbering each other. - * - * So: - * - **Reads** ([get], [containsKey], [size], [keys], [values]) are a single atomic load - * plus a lookup on an immutable map. No lock, no allocation, no copy. - * - **Bulk operations** (`filter`/`map`/`forEach`/…) iterate that same immutable map - * directly. No defensive `entries.toList()`, no `ConcurrentModificationException` - * window, and — because no lock is held while a caller's lambda runs — no way for a - * `LocalCache` predicate that reaches back into the cache to deadlock. - * - **Writes** are a CAS retry loop over a structurally shared copy. - * - * ## What it costs - * - * A write costs more than a `HashMap.put` under a lock would — a few node copies rather - * than one bucket store. Measured on this target (linuxX64, `-opt`, ms for the whole - * loop) against the copy-on-write version this replaces and against a - * `PlatformLock` + `HashMap` variant that was the other candidate: - * - * ``` - * n=20,000 fill reads 20 scans mixed (put + scan every 1k) - * copy-on-write 17,949 2 13 25,278 - * lock + HashMap 1 0 12 35 - * HAMT + CAS 14 0 19 28 - * - * n=200,000 fill reads 20 scans mixed - * lock + HashMap 44 9 177 6,736 - * HAMT + CAS 197 12 237 2,486 - * ``` - * - * Writing nothing but writes, the lock wins ~4x. But that is not the shape `LocalCache` - * has: it interleaves scans (every feed filters the whole cache) with arriving events, - * and there the lock has to rebuild an O(n) read snapshot after each write epoch, so it - * loses by ~2.7x at 200k. Reads and scans are close either way. The lock-free version - * therefore wins the workload that matters *and* is the one that never blocks a reader. - * - * This is the house pattern for shared mutable state in `commonMain` already — see - * `nip01Core.relay.filters.FilterIndex` and `nip86RelayManagement.server.BanStore`, - * both of which hold their state in one `AtomicReference` over persistent collections - * and mutate it with the same CAS loop. - * - * ## Notes - * - * Anything that reads twice — [remove], [getOrCreate], [createIfAbsent] — retries on a - * lost CAS rather than locking, so the pair is atomic without excluding readers. - * [clear] publishes an empty map unconditionally and can therefore drop a write that - * lands concurrently, exactly as `ConcurrentSkipListMap.clear()` can. - * - * Iteration order is hash order, as on Apple; JVM/Android is sorted-key order - * (`ConcurrentSkipListMap`). Nothing in the codebase depends on a specific order. The - * `from`/`to` range overloads below degrade to a full scan on this target, as they - * always have — they have no callers outside the JVM-only `LargeSoftCache`. + * Iteration is weakly consistent and in bucket order. JVM/Android iterates in + * sorted-key order (`ConcurrentSkipListMap`) and Apple in hash order; nothing in the + * codebase depends on a specific one. The `from`/`to` range overloads degrade to a full + * scan here, as they always have — they have no callers outside the JVM-only + * `LargeSoftCache`. */ -@OptIn(ExperimentalAtomicApi::class) actual class LargeCache : ICacheOperations { - private val ref = AtomicReference>(persistentHashMapOf()) + private val cache = StripedHashMap() - /** - * Runs [block] against an immutable point-in-time map. Outside any critical - * section, so [block] may call back into this cache freely. - */ - private inline fun withMap(block: (Map) -> R): R = block(ref.load()) - - /** - * Publishes [transform] of the current map with a CAS, retrying if a concurrent - * writer won. [transform] must be pure — it can run more than once — and returning - * the map it was given means "no change", which skips the CAS entirely. - */ - private inline fun mutate(transform: (PersistentMap) -> PersistentMap) { - while (true) { - val current = ref.load() - val next = transform(current) - if (next === current || ref.compareAndSet(current, next)) return - } + actual fun keys(): Set { + val results = LinkedHashSet(cache.size()) + cache.forEachEntry { key, _ -> results.add(key) } + return results } - actual fun keys(): Set = ref.load().keys - - actual fun values(): Iterable = ref.load().values - - actual fun get(key: K): V? = ref.load()[key] - - actual fun remove(key: K): V? { - while (true) { - val current = ref.load() - val previous = current[key] ?: return null - if (ref.compareAndSet(current, current.removing(key))) return previous - } + actual fun values(): Iterable { + val results = ArrayList(cache.size()) + cache.forEachEntry { _, value -> results.add(value) } + return results } - actual fun isEmpty(): Boolean = ref.load().isEmpty() + actual fun get(key: K): V? = cache.get(key) + + actual fun remove(key: K): V? = cache.remove(key) + + actual fun isEmpty(): Boolean = cache.isEmpty() actual fun clear() { - ref.store(persistentHashMapOf()) + cache.clear() } - actual fun containsKey(key: K): Boolean = ref.load().containsKey(key) + actual fun containsKey(key: K): Boolean = cache.containsKey(key) actual fun put( key: K, value: V, ) { - mutate { it.putting(key, value) } + cache.put(key, value) } - /** - * Mirrors the JVM actual's `putIfAbsent`: [builder] runs at most once — outside the - * retry loop, since it is caller code — and the value is only published if no one - * won the race in the meantime. - */ + // The next two are the JVM actual's bodies verbatim, over the same putIfAbsent + // contract: [builder] runs outside the write path, and the loser of a race keeps + // the winner's value. + actual fun getOrCreate( key: K, builder: (key: K) -> V, ): V { - ref.load()[key]?.let { return it } + val value = cache.get(key) - val newObject = builder(key) - - while (true) { - val current = ref.load() - current[key]?.let { return it } - if (ref.compareAndSet(current, current.putting(key, newObject))) return newObject + return if (value != null) { + value + } else { + val newObject = builder(key) + cache.putIfAbsent(key, newObject) ?: newObject } } /** - * True only when *this* call inserted the value — matching the JVM actual's - * `putIfAbsent(key, newObject) == null`. The previous implementation returned + * True only when *this* call inserted. An early implementation returned * `get(key) != null`, which also reported true when another thread had just created * the entry, double-firing whatever the caller does with a fresh key. */ @@ -181,121 +107,139 @@ actual class LargeCache : ICacheOperations { key: K, builder: (key: K) -> V, ): Boolean { - if (ref.load().containsKey(key)) return false - - val newObject = builder(key) - - while (true) { - val current = ref.load() - if (current.containsKey(key)) return false - if (ref.compareAndSet(current, current.putting(key, newObject))) return true + val value = cache.get(key) + return if (value != null) { + false + } else { + val newObject = builder(key) + cache.putIfAbsent(key, newObject) == null } } - actual override fun size(): Int = ref.load().size + actual override fun size(): Int = cache.size() actual override fun forEach(consumer: ICacheBiConsumer) { - // The map is immutable, so this iterates a stable snapshot with no copy. - withMap { map -> map.forEach { consumer.accept(it.key, it.value) } } + cache.forEachEntry { key, value -> consumer.accept(key, value) } } - actual override fun filter(consumer: CacheCollectors.BiFilter): List = withMap { map -> map.filter { consumer.filter(it.key, it.value) }.values.toList() } + actual override fun filter(consumer: CacheCollectors.BiFilter): List { + val results = ArrayList() + cache.forEachEntry { key, value -> if (consumer.filter(key, value)) results.add(value) } + return results + } - actual override fun filterIntoSet(consumer: CacheCollectors.BiFilter): Set = withMap { map -> map.filter { consumer.filter(it.key, it.value) }.values.toSet() } + actual override fun filterIntoSet(consumer: CacheCollectors.BiFilter): Set { + val results = LinkedHashSet() + cache.forEachEntry { key, value -> if (consumer.filter(key, value)) results.add(value) } + return results + } - actual override fun map(consumer: CacheCollectors.BiNotNullMapper): List = withMap { map -> map.map { consumer.map(it.key, it.value) } } + actual override fun map(consumer: CacheCollectors.BiNotNullMapper): List { + val results = ArrayList(cache.size()) + cache.forEachEntry { key, value -> results.add(consumer.map(key, value)) } + return results + } - actual override fun mapNotNull(consumer: CacheCollectors.BiMapper): List = withMap { map -> map.mapNotNull { consumer.map(it.key, it.value) } } + actual override fun mapNotNull(consumer: CacheCollectors.BiMapper): List { + val results = ArrayList() + cache.forEachEntry { key, value -> consumer.map(key, value)?.let { results.add(it) } } + return results + } - actual override fun mapNotNullIntoSet(consumer: CacheCollectors.BiMapper): Set = mapNotNull(consumer).toSet() + actual override fun mapNotNullIntoSet(consumer: CacheCollectors.BiMapper): Set { + val results = LinkedHashSet() + cache.forEachEntry { key, value -> consumer.map(key, value)?.let { results.add(it) } } + return results + } - actual override fun mapFlatten(consumer: CacheCollectors.BiMapper?>): List = withMap { map -> map.flatMap { entry -> consumer.map(entry.key, entry.value) ?: emptyList() } } + actual override fun mapFlatten(consumer: CacheCollectors.BiMapper?>): List { + val results = ArrayList() + cache.forEachEntry { key, value -> consumer.map(key, value)?.let { results.addAll(it) } } + return results + } - actual override fun mapFlattenIntoSet(consumer: CacheCollectors.BiMapper?>): Set = mapFlatten(consumer).toSet() + actual override fun mapFlattenIntoSet(consumer: CacheCollectors.BiMapper?>): Set { + val results = LinkedHashSet() + cache.forEachEntry { key, value -> consumer.map(key, value)?.let { results.addAll(it) } } + return results + } actual override fun maxOrNullOf( filter: CacheCollectors.BiFilter, comparator: Comparator, - ): V? = - withMap { map -> - var maxV: V? = null - map.forEach { - if (filter.filter(it.key, it.value)) { - if (maxV == null || comparator.compare(it.value, maxV) > 0) { - maxV = it.value - } + ): V? { + var maxV: V? = null + cache.forEachEntry { key, value -> + if (filter.filter(key, value)) { + if (maxV == null || comparator.compare(value, maxV) > 0) { + maxV = value } } - maxV } + return maxV + } - actual override fun sumOf(consumer: CacheCollectors.BiSumOf): Int = - withMap { map -> - var sum = 0 - map.forEach { sum += consumer.map(it.key, it.value) } - sum - } + actual override fun sumOf(consumer: CacheCollectors.BiSumOf): Int { + var sum = 0 + cache.forEachEntry { key, value -> sum += consumer.map(key, value) } + return sum + } - actual override fun sumOfLong(consumer: CacheCollectors.BiSumOfLong): Long = - withMap { map -> - var sum = 0L - map.forEach { sum += consumer.map(it.key, it.value) } - sum - } + actual override fun sumOfLong(consumer: CacheCollectors.BiSumOfLong): Long { + var sum = 0L + cache.forEachEntry { key, value -> sum += consumer.map(key, value) } + return sum + } - actual override fun groupBy(consumer: CacheCollectors.BiNotNullMapper): Map> = - withMap { map -> - val results = HashMap>() - map.forEach { - val group = consumer.map(it.key, it.value) - results.getOrPut(group) { ArrayList() }.add(it.value) - } - results + actual override fun groupBy(consumer: CacheCollectors.BiNotNullMapper): Map> { + val results = HashMap>() + cache.forEachEntry { key, value -> + results.getOrPut(consumer.map(key, value)) { ArrayList() }.add(value) } + return results + } - actual override fun countByGroup(consumer: CacheCollectors.BiNotNullMapper): Map = - withMap { map -> - val results = HashMap() - map.forEach { - val group = consumer.map(it.key, it.value) - results[group] = (results[group] ?: 0) + 1 - } - results + actual override fun countByGroup(consumer: CacheCollectors.BiNotNullMapper): Map { + val results = HashMap() + cache.forEachEntry { key, value -> + val group = consumer.map(key, value) + results[group] = (results[group] ?: 0) + 1 } + return results + } actual override fun sumByGroup( groupMap: CacheCollectors.BiNotNullMapper, sumOf: CacheCollectors.BiNotNullMapper, - ): Map = - withMap { map -> - val results = HashMap() - map.forEach { - val group = groupMap.map(it.key, it.value) - results[group] = (results[group] ?: 0L) + sumOf.map(it.key, it.value) - } - results + ): Map { + val results = HashMap() + cache.forEachEntry { key, value -> + val group = groupMap.map(key, value) + results[group] = (results[group] ?: 0L) + sumOf.map(key, value) } + return results + } - actual override fun count(consumer: CacheCollectors.BiFilter): Int = withMap { map -> map.count { consumer.filter(it.key, it.value) } } + actual override fun count(consumer: CacheCollectors.BiFilter): Int { + var count = 0 + cache.forEachEntry { key, value -> if (consumer.filter(key, value)) count++ } + return count + } - actual override fun associate(transform: (K, V) -> Pair): Map = - withMap { map -> - val results = LinkedHashMap(map.size) - map.forEach { - val pair = transform(it.key, it.value) - results[pair.first] = pair.second - } - results + actual override fun associate(transform: (K, V) -> Pair): Map { + val results = LinkedHashMap(cache.size()) + cache.forEachEntry { key, value -> + val pair = transform(key, value) + results[pair.first] = pair.second } + return results + } - actual override fun associateWith(transform: (K, V) -> U?): Map = - withMap { map -> - val results = LinkedHashMap(map.size) - map.forEach { - results[it.key] = transform(it.key, it.value) - } - results - } + actual override fun associateWith(transform: (K, V) -> U?): Map { + val results = LinkedHashMap(cache.size()) + cache.forEachEntry { key, value -> results[key] = transform(key, value) } + return results + } actual override fun filter( from: K, diff --git a/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/StripedHashMap.kt b/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/StripedHashMap.kt new file mode 100644 index 0000000000..ec9caba654 --- /dev/null +++ b/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/StripedHashMap.kt @@ -0,0 +1,309 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.quartz.utils.cache + +import com.vitorpamplona.quartz.utils.concurrent.PlatformLock +import com.vitorpamplona.quartz.utils.concurrent.withLock +import kotlin.concurrent.Volatile +import kotlin.concurrent.atomics.AtomicArray +import kotlin.concurrent.atomics.AtomicInt +import kotlin.concurrent.atomics.ExperimentalAtomicApi + +/** + * A chained hash table with **lock-free reads** and **striped-lock writes** — the shape + * of `java.util.concurrent.ConcurrentHashMap`, which Kotlin/Native has no equivalent of. + * + * Exists because `LocalCache` fills on the order of 100,000 entries in a few seconds, + * and every structure available on this target fails that workload in some way: + * + * - **Copy-on-write over a `HashMap`** (what shipped first) rebuilds the whole map per + * write: O(n) each, O(n^2) to fill. + * - **A persistent HAMT + CAS** is O(log32 n) per write, but allocates a fresh path of + * ~4-5 nodes for *every* write — including overwrites, which change no structure at + * all — and throws the old path away. Measured over a 100k fill plus scans that is 24 + * GC cycles against this table's 1. + * - **One lock around a `HashMap`** writes fast but has to hand bulk operations an O(n) + * copy, because a caller's lambda must not run inside the critical section (the + * linux `PlatformLock` is a spin lock and is not reentrant, and `LocalCache` + * predicates reach back into the cache). + * + * A chained table avoids all three. Structure is only touched when a key is *added* + * (one node, prepended), an overwrite is a single volatile store into the existing + * node, and scans walk the buckets in place with no copy and no lock — so caller + * lambdas run outside any critical section and cannot deadlock. + * + * Measured on linuxX64 (`-opt`), 100,000 String keys of event-id length, ms per phase + * and GC cycles over the whole run: + * + * ``` + * fill overwrite reads 20 scans mixed GCs heap + * copy-on-write* n/a n/a n/a n/a n/a n/a n/a + * HAMT + CAS 70 78 6 117 676 24 67MB + * lock + HashMap 13 6 3 71 1197 36 51MB + * this 15 3 3 13 64 1 43MB + * ``` + * + * (*copy-on-write is off the scale: 20k entries alone took 18s to fill.) "mixed" is a + * full fill with a whole-table scan every 1000 writes, which is the shape `LocalCache` + * actually has — arriving events interleaved with feeds filtering the whole cache. + * Figures are one representative run of several; they were stable to within ~10%, + * except the HAMT's scan column, which wandered between 115ms and 190ms. + * + * ## Concurrency contract + * + * - **Readers never block and never allocate.** [get], [containsKey], [size] and + * [forEachEntry] take no lock. A reader loads [table] once and walks immutable + * `next` links, so it always sees a well-formed chain. + * - **Writers block only against writers hashing to the same stripe**, and only for a + * bucket walk of a few nodes. This is where it differs from the JVM actual's fully + * non-blocking `ConcurrentSkipListMap`; it matches `ConcurrentHashMap`, which also + * locks a bin to write it. + * - Iteration is **weakly consistent**, like both of those: it reflects the table as of + * its first load and may or may not observe writes that land while it runs. It never + * throws, never sees a torn chain, and never needs a defensive copy. + * - A resize takes every stripe lock, so no write can be in flight while it runs. + * Nodes are rebuilt rather than relinked, which is what lets a reader that captured + * the pre-resize table keep walking it safely. + */ +@OptIn(ExperimentalAtomicApi::class) +internal class StripedHashMap { + /** + * [value] is mutable so that overwriting an existing key allocates nothing; [next] + * is not, so that a reader walking a chain can never see it change under them. + * Structural edits publish a new head instead. + */ + internal class Node( + val hash: Int, + val key: K, + @Volatile var value: V, + val next: Node?, + ) + + private val locks = Array(STRIPES) { PlatformLock() } + private val entryCount = AtomicInt(0) + + @Volatile internal var table = AtomicArray?>(INITIAL_CAPACITY) { null } + + @Volatile private var threshold = INITIAL_CAPACITY / 4 * 3 + + /** Spreads the high bits down, so that both the bucket and the stripe see entropy. */ + private fun hashOf(key: K): Int { + val h = key?.hashCode() ?: 0 + return h xor (h ushr 16) + } + + /** + * Derived from the hash alone, never from the table size, so a key keeps the same + * stripe across a resize. + */ + private fun lockFor(hash: Int) = locks[(hash ushr 16) and (STRIPES - 1)] + + fun size(): Int = entryCount.load() + + fun isEmpty(): Boolean = entryCount.load() == 0 + + fun get(key: K): V? { + val hash = hashOf(key) + val current = table + var node = current.loadAt(hash and (current.size - 1)) + while (node != null) { + if (node.hash == hash && node.key == key) return node.value + node = node.next + } + return null + } + + fun containsKey(key: K): Boolean { + val hash = hashOf(key) + val current = table + var node = current.loadAt(hash and (current.size - 1)) + while (node != null) { + if (node.hash == hash && node.key == key) return true + node = node.next + } + return false + } + + fun put( + key: K, + value: V, + ) { + val hash = hashOf(key) + var grew = false + lockFor(hash).withLock { + val current = table + val index = hash and (current.size - 1) + val head = current.loadAt(index) + var node = head + while (node != null) { + if (node.hash == hash && node.key == key) { + // Present already: no structural change, no allocation. + node.value = value + return@withLock + } + node = node.next + } + current.storeAt(index, Node(hash, key, value, head)) + grew = entryCount.fetchAndAdd(1) + 1 > threshold + } + if (grew) growTable() + } + + /** + * Inserts [value] only if [key] is absent, and returns the value already stored — + * or null when this call performed the insert. Exactly `ConcurrentMap.putIfAbsent`, + * which the JVM actual builds `getOrCreate` and `createIfAbsent` out of, including + * its inability to represent a stored null (`ConcurrentSkipListMap` rejects those). + */ + fun putIfAbsent( + key: K, + value: V, + ): V? { + val hash = hashOf(key) + var grew = false + var existing: V? = null + lockFor(hash).withLock { + val current = table + val index = hash and (current.size - 1) + val head = current.loadAt(index) + var node = head + while (node != null) { + if (node.hash == hash && node.key == key) { + existing = node.value + return@withLock + } + node = node.next + } + current.storeAt(index, Node(hash, key, value, head)) + grew = entryCount.fetchAndAdd(1) + 1 > threshold + } + if (grew) growTable() + return existing + } + + fun remove(key: K): V? { + val hash = hashOf(key) + var removed: V? = null + lockFor(hash).withLock { + val current = table + val index = hash and (current.size - 1) + val head = current.loadAt(index) + + var target = head + while (target != null && !(target.hash == hash && target.key == key)) target = target.next + if (target == null) return@withLock + + // `next` is immutable, so the nodes ahead of the removed one are cloned onto + // its tail. A reader still walking the old head sees the entry one last time + // rather than a broken chain. + var rebuilt = target.next + var ahead = head + while (ahead !== target) { + val node = ahead!! + rebuilt = Node(node.hash, node.key, node.value, rebuilt) + ahead = node.next + } + + current.storeAt(index, rebuilt) + entryCount.fetchAndAdd(-1) + removed = target.value + } + return removed + } + + fun clear() { + lockAll() + try { + table = AtomicArray(INITIAL_CAPACITY) { null } + threshold = INITIAL_CAPACITY / 4 * 3 + entryCount.store(0) + } finally { + unlockAll() + } + } + + /** + * Walks every entry without locking. Inline so the caller's body runs with no + * `Function2` dispatch and no captured-variable box per entry, which is what keeps + * a full-cache scan allocation-free. + */ + inline fun forEachEntry(action: (K, V) -> Unit) { + val current = table + for (index in 0 until current.size) { + var node = current.loadAt(index) + while (node != null) { + action(node.key, node.value) + node = node.next + } + } + } + + private fun growTable() { + lockAll() + try { + val old = table + // Another writer may have grown it while this one waited for the locks. + if (entryCount.load() <= threshold) return + if (old.size >= MAX_CAPACITY) { + threshold = Int.MAX_VALUE + return + } + + val capacity = old.size shl 1 + val next = AtomicArray?>(capacity) { null } + for (index in 0 until old.size) { + var node = old.loadAt(index) + while (node != null) { + val target = node.hash and (capacity - 1) + next.storeAt(target, Node(node.hash, node.key, node.value, next.loadAt(target))) + node = node.next + } + } + + table = next + threshold = capacity / 4 * 3 + } finally { + unlockAll() + } + } + + /** Always in index order, and only ever from a thread holding no stripe lock. */ + private fun lockAll() { + for (lock in locks) lock.lock() + } + + private fun unlockAll() { + for (index in locks.indices.reversed()) locks[index].unlock() + } + + companion object { + /** + * Writes are O(1), so a stripe is held for a few nanoseconds and 16 ways is + * plenty — `ConcurrentHashMap` shipped with the same default for years. + */ + private const val STRIPES = 16 + + /** Sized to carry a warm cache's first few thousand entries without a resize. */ + private const val INITIAL_CAPACITY = 1024 + + private const val MAX_CAPACITY = 1 shl 30 + } +} diff --git a/quartz/src/linuxTest/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCacheCollisionTest.kt b/quartz/src/linuxTest/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCacheCollisionTest.kt new file mode 100644 index 0000000000..b4ccaeca98 --- /dev/null +++ b/quartz/src/linuxTest/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCacheCollisionTest.kt @@ -0,0 +1,140 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.quartz.utils.cache + +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertFalse +import kotlin.test.assertNull +import kotlin.test.assertTrue + +/** + * Drives every key into a single bucket of the [StripedHashMap] backing this target's + * [LargeCache], so the bucket-chain paths run deterministically instead of only when a + * hash happens to collide. + * + * The one that most needs it is removal. Chain nodes hold their `next` immutably — that + * is what lets a reader walk a chain with no lock — so removing from the middle has to + * clone the nodes ahead of the target onto its tail and publish a new head. With + * well-spread keys that path almost never sees a chain longer than two. + */ +class LargeCacheCollisionTest { + /** Every instance lands in the same bucket, and in the same stripe. */ + private data class Collides( + val id: Int, + ) { + override fun hashCode() = 0 + } + + private fun filled(n: Int) = + LargeCache().apply { + for (i in 0 until n) put(Collides(i), i) + } + + @Test + fun readsFindEveryEntryInOneChain() { + val cache = filled(200) + + assertEquals(200, cache.size()) + for (i in 0 until 200) { + assertEquals(i, cache.get(Collides(i)), "entry $i") + assertTrue(cache.containsKey(Collides(i))) + } + assertNull(cache.get(Collides(200))) + assertFalse(cache.containsKey(Collides(200))) + } + + @Test + fun overwriteInAChainReplacesInPlace() { + val cache = filled(200) + + for (i in 0 until 200) cache.put(Collides(i), i * 10) + + assertEquals(200, cache.size(), "overwriting must not lengthen the chain") + for (i in 0 until 200) assertEquals(i * 10, cache.get(Collides(i))) + } + + @Test + fun removeFromTheMiddleKeepsTheRestOfTheChain() { + val cache = filled(200) + + // Head, tail and middle of the chain, in an order that leaves gaps behind. + for (i in 0 until 200 step 3) { + assertEquals(i, cache.remove(Collides(i)), "remove $i returns its value") + } + + val expected = (0 until 200).filter { it % 3 != 0 } + assertEquals(expected.size, cache.size()) + for (i in 0 until 200) { + if (i % 3 == 0) { + assertNull(cache.get(Collides(i)), "entry $i was removed") + } else { + assertEquals(i, cache.get(Collides(i)), "entry $i survived") + } + } + + val seen = mutableListOf() + cache.forEach { _, v -> seen.add(v) } + assertEquals(expected.toSet(), seen.toSet(), "iteration must match the survivors") + assertEquals(expected.size, seen.size, "iteration must not double-count") + + assertNull(cache.remove(Collides(0)), "removing twice is a no-op") + assertEquals(expected.size, cache.size()) + } + + @Test + fun getOrCreateAndCreateIfAbsentWalkTheChain() { + val cache = filled(200) + var builds = 0 + + for (i in 0 until 200) { + assertEquals( + i, + cache.getOrCreate(Collides(i)) { + builds++ + -1 + }, + ) + assertFalse(cache.createIfAbsent(Collides(i)) { -1 }) + } + assertEquals(0, builds, "nothing in the chain should have been rebuilt") + + assertTrue(cache.createIfAbsent(Collides(500)) { 500 }) + assertEquals(500, cache.get(Collides(500))) + assertEquals(201, cache.size()) + } + + @Test + fun clearEmptiesAFullBucket() { + val cache = filled(200) + + cache.clear() + + assertEquals(0, cache.size()) + assertTrue(cache.isEmpty()) + assertNull(cache.get(Collides(7))) + assertEquals(0, cache.count { _, _ -> true }) + + cache.put(Collides(1), 1) + assertEquals(1, cache.size()) + assertEquals(1, cache.get(Collides(1))) + } +} From d759741f9641c7909eb186a1201d17fce2bc5738 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 1 Sep 2026 19:06:53 +0000 Subject: [PATCH 4/6] fix(quartz): make the whole linuxX64 test suite pass; widen the CI leg MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The 78 failures on this target were not 78 unimplemented actuals. Two root causes accounted for all of them. TestResourceLoader.linux was a TODO(), so every vector-driven suite failed before it reached any production code: the full MLS interop set, NIP-44, the NIP-01 hint indexer, the SQLite store's large-DB tests and the Bolt12 payer proofs — 69 tests. Implemented over platform.posix (linuxX64 has no Foundation for the Apple actual's NSData path), resolving against the same TEST_RESOURCES_ROOT that build.gradle.kts already exports onto every KotlinNativeTest task. The read is one ftell-sized allocation filled by fread, so a vector file costs exactly one ByteArray — less than the JVM actual's bufferedReader().readText(), which grows a StringBuilder as it goes. UriParser.linux never URL-decoded query values or fragments, though the JVM actual runs both through URLDecoder.decode(.., "UTF-8"). Every NIP-47 failure was one symptom of that: relay=wss%3A%2F%2Frelay.damus.io reached RelayUrlNormalizer still percent-encoded and came back "Invalid relay Url" (6 tests), and the deep-link round trips compared an encoded string against a plain one (3 tests). Added a decoder matching URLDecoder where the behaviour is observable — '+' to space, a run of consecutive %XX decoded as one UTF-8 sequence, malformed escapes throwing IllegalArgumentException — with the same short-circuit URLDecoder makes, returning the original instance when there is nothing to decode. Two other divergences fixed while there: getQueryParameter returned an empty list where the JVM returns null for an absent parameter, and the query string was re-split on every call rather than parsed once into a lazy map, so a URI read for four parameters was parsed four times. With those, linuxX64Test is 3495 tests, 0 failures, so the CI leg added alongside the LargeCache work drops its cache-package filter and runs the whole :quartz suite. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01HxQ1QuyzSkR38iFHbREjoS --- .github/workflows/build.yml | 32 ++--- .../quartz/utils/UriParser.linux.kt | 134 +++++++++++++++--- .../quartz/TestResourceLoader.linux.kt | 82 ++++++++++- 3 files changed, 197 insertions(+), 51 deletions(-) diff --git a/.github/workflows/build.yml b/.github/workflows/build.yml index c130551bce..7c55e0819b 100644 --- a/.github/workflows/build.yml +++ b/.github/workflows/build.yml @@ -194,24 +194,15 @@ jobs: name: geode Test Reports path: geode/build/reports - # linuxX64 is the only target whose LargeCache / ConcurrentHashCache actuals are - # hand-written concurrent maps rather than a delegation to a platform concurrent - # collection, and until this job existed nothing ran them: the target was compiled - # by no CI leg at all. That is how a copy-on-write LargeCache with O(n) writes and a - # non-atomic read-copy-write (concurrent writers silently dropped entries) sat in - # the tree unnoticed. + # Until this job existed nothing ran the linuxX64 target at all — it was compiled by + # no CI leg. That is how a copy-on-write LargeCache with O(n) writes and a non-atomic + # read-copy-write (concurrent writers silently dropped entries) sat in the tree + # unnoticed, and how TestResourceLoader stayed a TODO() that failed every vector-driven + # suite on the target. # - # Scoped to the cache/concurrency packages on purpose. The full linuxX64Test suite - # is 3,490 tests with 78 pre-existing failures, essentially all of them `TODO()` - # stubs in linux actuals that were never written — MLS crypto, the SQLite driver, - # NIP-44, Bolt12 — plus two URL-handling divergences. Filling those in is its own - # project; gating PRs on them today would just mean a permanently red job. The - # filter keeps the leg meaningful and green, and widening it is a one-line change - # once the native actuals land. - # - # This still compiles and links the whole module for linuxX64, so a commonMain or - # commonTest source that reaches for a JVM-only API fails here too — on a target - # with no Foundation to fall back on the way Apple has. + # Runs the whole :quartz suite on a Linux Native frontend, which also catches a + # commonMain or commonTest source reaching for a JVM-only API on a target that, unlike + # Apple, has no Foundation to fall back on. test-quartz-linux-native: needs: lint runs-on: ubuntu-latest @@ -242,11 +233,8 @@ jobs: key: konan-${{ runner.os }}-${{ hashFiles('gradle/libs.versions.toml') }} restore-keys: konan-${{ runner.os }}- - - name: Test Quartz caches on Linux Native - run: | - ./gradlew :quartz:linuxX64Test \ - --tests "com.vitorpamplona.quartz.utils.cache.*" \ - --tests "com.vitorpamplona.quartz.utils.concurrent.*" + - name: Test Quartz on Linux Native + run: ./gradlew :quartz:linuxX64Test - name: Linux Native Test Report uses: mikepenz/action-junit-report@a9170d5795813c01ab4901ffb045b52bab4ab09d # v6.5.0 diff --git a/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/UriParser.linux.kt b/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/UriParser.linux.kt index 5f5d9f0c27..448f4b8c8f 100644 --- a/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/UriParser.linux.kt +++ b/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/UriParser.linux.kt @@ -121,41 +121,129 @@ actual class UriParser actual constructor( actual fun path(): String? = parsedPath - actual fun queryParameterNames(): Set { - val query = parsedQuery ?: return emptySet() - return query - .split('&') - .map { param -> - val eqIndex = param.indexOf('=') - if (eqIndex >= 0) param.substring(0, eqIndex) else param - }.toSet() - } + /** + * Parsed once and reused, mirroring the JVM actual's lazy map. The previous version + * re-split the entire query string on every [getQueryParameter] call, so a URI read + * for four parameters was parsed four times. + */ + private val queryParameters: Map> by lazy { + parsedQuery?.ifBlank { null }?.let { query -> + val params = mutableMapOf>() - actual fun getQueryParameter(param: String): List? { - val query = parsedQuery ?: return null - return query - .split('&') - .filter { part -> - val eqIndex = part.indexOf('=') - if (eqIndex >= 0) part.substring(0, eqIndex) == param else part == param - }.map { part -> - val eqIndex = part.indexOf('=') - if (eqIndex >= 0) part.substring(eqIndex + 1) else "" + query.split('&').forEach { paramValue -> + val parts = paramValue.split("=", limit = 2) + val currentValue = + params.getOrPut(parts[0]) { + mutableListOf() + } + + if (parts.size == 2) { + currentValue.add(percentDecode(parts[1])) + } else { + currentValue.add("") + } } + + params + } ?: emptyMap() } - val fragments: Map by lazy { + private val parsedFragments: Map by lazy { parsedFragment?.ifBlank { null }?.let { keyValuePair -> keyValuePair.split('&').associate { paramValue -> val parts = paramValue.split("=", limit = 2) if (parts.size == 2) { - parts[0] to parts[1] + parts[0] to percentDecode(parts[1]) } else { - parts[0] to "" + parts[0] to "" // Handle parameters without a value } } } ?: emptyMap() } - actual fun fragments(): Map = fragments + actual fun queryParameterNames(): Set = queryParameters.keys + + /** Null — not an empty list — when the parameter is absent, as on the JVM. */ + actual fun getQueryParameter(param: String): List? = queryParameters[param] + + actual fun fragments(): Map = parsedFragments } + +/** + * `java.net.URLDecoder.decode(value, "UTF-8")` — which is literally what the JVM actual + * calls — for a target with no `java.net`. + * + * This is the whole reason NIP-47 failed on linuxX64: the parser returned query values + * exactly as they appeared in the URI, so `relay=wss%3A%2F%2Frelay.damus.io` reached + * `RelayUrlNormalizer` still percent-encoded and came back "Invalid relay Url". Decoding + * belongs here rather than in each caller, because the JVM and Apple actuals both hand + * back decoded values and common code is written against that. + * + * Matches `URLDecoder` in the details that are observable: `+` becomes a space, a run of + * consecutive `%XX` is decoded as one UTF-8 sequence (so multi-byte characters survive), + * every other character passes through, and a malformed escape throws + * [IllegalArgumentException] rather than being silently kept — the same failure the JVM + * gives for the same input. + */ +private fun percentDecode(value: String): String { + // The short-circuit URLDecoder also makes: with nothing to change, return the + // original instance rather than rebuilding it. Most query values hit this. + if (value.indexOf('%') < 0 && value.indexOf('+') < 0) return value + + val result = StringBuilder(value.length) + var index = 0 + // Sized on first use for the longest run that could still follow, then reused — + // one allocation for the whole string, as in URLDecoder. + var escaped: ByteArray? = null + + while (index < value.length) { + when (val char = value[index]) { + '+' -> { + result.append(' ') + index++ + } + + '%' -> { + val buffer = escaped ?: ByteArray((value.length - index) / 3).also { escaped = it } + var count = 0 + while (index + 2 < value.length && value[index] == '%') { + buffer[count++] = decodeEscape(value, index) + index += 3 + } + if (index < value.length && value[index] == '%') { + throw IllegalArgumentException("URLDecoder: Incomplete trailing escape (%) pattern") + } + result.append(buffer.decodeToString(0, count)) + } + + else -> { + result.append(char) + index++ + } + } + } + + return result.toString() +} + +private fun decodeEscape( + value: String, + index: Int, +): Byte { + val high = hexDigit(value[index + 1]) + val low = hexDigit(value[index + 2]) + if (high < 0 || low < 0) { + throw IllegalArgumentException( + "URLDecoder: Illegal hex characters in escape (%) pattern - ${value.substring(index, index + 3)}", + ) + } + return ((high shl 4) or low).toByte() +} + +private fun hexDigit(char: Char): Int = + when (char) { + in '0'..'9' -> char - '0' + in 'a'..'f' -> char - 'a' + 10 + in 'A'..'F' -> char - 'A' + 10 + else -> -1 + } diff --git a/quartz/src/linuxTest/kotlin/com/vitorpamplona/quartz/TestResourceLoader.linux.kt b/quartz/src/linuxTest/kotlin/com/vitorpamplona/quartz/TestResourceLoader.linux.kt index 3997dd8575..3d8f191d19 100644 --- a/quartz/src/linuxTest/kotlin/com/vitorpamplona/quartz/TestResourceLoader.linux.kt +++ b/quartz/src/linuxTest/kotlin/com/vitorpamplona/quartz/TestResourceLoader.linux.kt @@ -20,12 +20,82 @@ */ package com.vitorpamplona.quartz -actual class TestResourceLoader actual constructor() { - actual fun loadDecompressString(file: String): String { - TODO("Not yet implemented") - } +import com.vitorpamplona.quartz.utils.GZip +import kotlinx.cinterop.ExperimentalForeignApi +import kotlinx.cinterop.addressOf +import kotlinx.cinterop.convert +import kotlinx.cinterop.toKString +import kotlinx.cinterop.usePinned +import platform.posix.SEEK_END +import platform.posix.SEEK_SET +import platform.posix.fclose +import platform.posix.fopen +import platform.posix.fread +import platform.posix.fseek +import platform.posix.ftell +import platform.posix.getenv - actual fun loadString(file: String): String { - TODO("Not yet implemented") +/** + * Linux/Native actual for [TestResourceLoader]. + * + * Until this existed it was `TODO()`, which failed 69 tests on this target — every + * suite driven by a vector file: the whole MLS interop set, NIP-44, the NIP-01 hint + * indexer, the SQLite store's large-DB tests and the Bolt12 payer proofs. None of them + * were failing because of missing production code; they could not read their input. + * + * Resolves paths against `TEST_RESOURCES_ROOT`, the same environment variable the Apple + * actual uses, exported onto every `KotlinNativeTest` task by `quartz/build.gradle.kts`. + * Reads through `platform.posix` rather than Foundation, which linuxX64 does not have. + * + * The read is a single `stat`-sized allocation filled by `fread`, so a vector file + * costs exactly one `ByteArray` — less than the JVM actual's `bufferedReader().readText()`, + * which grows a `StringBuilder` as it goes. + */ +@OptIn(ExperimentalForeignApi::class) +actual class TestResourceLoader actual constructor() { + actual fun loadDecompressString(file: String): String = GZip.decompress(readBytes(file)) + + actual fun loadString(file: String): String = readBytes(file).decodeToString() + + private fun readBytes(file: String): ByteArray { + val root = + getenv("TEST_RESOURCES_ROOT")?.toKString() + ?: throw IllegalStateException( + "TEST_RESOURCES_ROOT is not set. quartz/build.gradle.kts exports it onto every " + + "KotlinNativeTest task; running the test binary directly has to set it too.", + ) + + val path = "$root/$file" + val handle = fopen(path, "rb") ?: throw IllegalArgumentException("Resource not found: $path") + + try { + if (fseek(handle, 0, SEEK_END) != 0) throw IllegalArgumentException("Resource is not seekable: $path") + val size = ftell(handle) + if (size < 0L) throw IllegalArgumentException("Cannot determine the size of: $path") + if (size == 0L) return ByteArray(0) + if (fseek(handle, 0, SEEK_SET) != 0) throw IllegalArgumentException("Cannot rewind: $path") + + val bytes = ByteArray(size.toInt()) + bytes.usePinned { pinned -> + var read = 0 + while (read < bytes.size) { + val count = + fread( + pinned.addressOf(read), + 1.convert(), + (bytes.size - read).convert(), + handle, + ).toInt() + if (count <= 0) break + read += count + } + if (read != bytes.size) { + throw IllegalArgumentException("Short read on $path: got $read of ${bytes.size} bytes") + } + } + return bytes + } finally { + fclose(handle) + } } } From fa3287f73719e64cda70983d35ebf311b4e51add Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 1 Sep 2026 19:38:00 +0000 Subject: [PATCH 5/6] fix(quartz): match java.net.URLEncoder on native; drop the urlencoder dep MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Both native targets delegated UrlEncoder to net.thauvin.erik.urlencoder.UrlEncoderUtil, which implements RFC 3986 percent-encoding. The JVM/Android actual is java.net.URLEncoder/URLDecoder, which implements application/x-www-form-urlencoded. Different specifications, and the difference was observable: JVM/Android UrlEncoderUtil encode(" ") "+" "%20" encode("*") "*" "%2A" decode("a+b") "a b" "a+b" This is not cosmetic. encode() builds strings that leave the device — TorrentEvent puts it in magnet links, Nip54InlineMetadata in inline metadata, Nip47DeepLink in the callback/appname/value parameters of NWC deep links — so Android and iOS emitted different bytes for the same title. The decode row is worse: a link written by Android carries '+' for its spaces, and reading it on iOS or desktop-native gave back literal plus signs, silently, with no error. Replaced with one UrlEncoder.native.kt in nativeMain, shared by linuxX64 and every Apple target, matching URLEncoder/URLDecoder exactly — unreserved set is alphanumerics plus -_.* (note '*' survives and '~' does not, the opposite of RFC 3986), space to '+', uppercase %XX of UTF-8 bytes otherwise, and '+' back to space on the way in. Escape runs are encoded and decoded as runs so surrogate pairs and multi-byte sequences survive, and both directions short-circuit on a string with nothing to change, as the java.net pair does. UriParser.linux now delegates to UrlEncoder.decode rather than carrying its own copy of the decoder added in the previous commit. The new UrlEncoderTest lives in commonTest, so it pins every target against the JVM's answers — it is what found all three rows above, by passing on jvmTest and failing three of ten on linuxX64. net.thauvin.erik:urlencoder-lib had no other user and is removed from both source sets and the version catalog. One deliberate edge difference from the JVM, documented at the call site: an unpaired UTF-16 surrogate encodes as %EF%BF%BD rather than %3F. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01HxQ1QuyzSkR38iFHbREjoS --- gradle/libs.versions.toml | 2 - quartz/build.gradle.kts | 2 - .../quartz/utils/UrlEncoder.apple.kt | 29 --- .../quartz/utils/UrlEncoderTest.kt | 134 ++++++++++++ .../quartz/utils/UriParser.linux.kt | 93 +-------- .../quartz/utils/UrlEncoder.linux.kt | 29 --- .../quartz/utils/UrlEncoder.native.kt | 190 ++++++++++++++++++ 7 files changed, 333 insertions(+), 146 deletions(-) delete mode 100644 quartz/src/appleMain/kotlin/com/vitorpamplona/quartz/utils/UrlEncoder.apple.kt create mode 100644 quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/utils/UrlEncoderTest.kt delete mode 100644 quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/UrlEncoder.linux.kt create mode 100644 quartz/src/nativeMain/kotlin/com/vitorpamplona/quartz/utils/UrlEncoder.native.kt diff --git a/gradle/libs.versions.toml b/gradle/libs.versions.toml index f19b09c9ca..baeabf40c1 100644 --- a/gradle/libs.versions.toml +++ b/gradle/libs.versions.toml @@ -55,7 +55,6 @@ media3 = "1.11.0" mockk = "1.14.11" kotlinx-coroutines-test = "1.11.0" negentropyKmp = "v1.2.0" -netUrlencoderLibVersion = "1.6.0" navigationCompose = "2.9.8" okhttp = "5.5.0" osmdroid = "6.1.20" @@ -209,7 +208,6 @@ mockk = { group = "io.mockk", name = "mockk", version.ref = "mockk" } mockk-android = { group = "io.mockk", name = "mockk-android", version.ref = "mockk" } kotlinx-coroutines-test = { group = "org.jetbrains.kotlinx", name = "kotlinx-coroutines-test", version.ref = "kotlinx-coroutines-test"} negentropy-kmp = { module = "com.vitorpamplona.negentropy:kmp-negentropy", version.ref = "negentropyKmp" } -net-thauvin-erik-urlencoder-lib = { module = "net.thauvin.erik.urlencoder:urlencoder-lib", version.ref = "netUrlencoderLibVersion" } okhttp = { group = "com.squareup.okhttp3", name = "okhttp", version.ref = "okhttp" } okhttpCoroutines = { group = "com.squareup.okhttp3", name = "okhttp-coroutines", version.ref = "okhttp" } osmdroid-android = { group = "org.osmdroid", name = "osmdroid-android", version.ref = "osmdroid" } diff --git a/quartz/build.gradle.kts b/quartz/build.gradle.kts index 4771e49055..9253f2c607 100644 --- a/quartz/build.gradle.kts +++ b/quartz/build.gradle.kts @@ -289,7 +289,6 @@ kotlin { dependsOn(nativeMain) dependencies { implementation(libs.charlietap.cachemap) - implementation(libs.net.thauvin.erik.urlencoder.lib) implementation(libs.dev.whyoleg.cryptography.provider.apple.optimal) implementation("io.github.andreypfau:kotlinx-crypto-hmac:0.0.4") implementation("io.github.andreypfau:kotlinx-crypto-sha2:0.0.4") @@ -347,7 +346,6 @@ kotlin { create("linuxMain") { dependsOn(nativeMain) dependencies { - implementation(libs.net.thauvin.erik.urlencoder.lib) implementation(libs.dev.whyoleg.cryptography.provider.apple.optimal) } } diff --git a/quartz/src/appleMain/kotlin/com/vitorpamplona/quartz/utils/UrlEncoder.apple.kt b/quartz/src/appleMain/kotlin/com/vitorpamplona/quartz/utils/UrlEncoder.apple.kt deleted file mode 100644 index 3b8e98f026..0000000000 --- a/quartz/src/appleMain/kotlin/com/vitorpamplona/quartz/utils/UrlEncoder.apple.kt +++ /dev/null @@ -1,29 +0,0 @@ -/* - * Copyright (c) 2025 Vitor Pamplona - * - * Permission is hereby granted, free of charge, to any person obtaining a copy of - * this software and associated documentation files (the "Software"), to deal in - * the Software without restriction, including without limitation the rights to use, - * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the - * Software, and to permit persons to whom the Software is furnished to do so, - * subject to the following conditions: - * - * The above copyright notice and this permission notice shall be included in all - * copies or substantial portions of the Software. - * - * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR - * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS - * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR - * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN - * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION - * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. - */ -package com.vitorpamplona.quartz.utils - -import net.thauvin.erik.urlencoder.UrlEncoderUtil - -actual object UrlEncoder { - actual fun encode(value: String): String = UrlEncoderUtil.encode(value) - - actual fun decode(value: String): String = UrlEncoderUtil.decode(value) -} diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/utils/UrlEncoderTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/utils/UrlEncoderTest.kt new file mode 100644 index 0000000000..33dc069b55 --- /dev/null +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/utils/UrlEncoderTest.kt @@ -0,0 +1,134 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.quartz.utils + +import kotlin.test.Test +import kotlin.test.assertEquals + +/** + * Cross-target contract for [UrlEncoder]. + * + * This is not a style preference — [UrlEncoder.encode] builds strings that leave the + * device. `TorrentEvent` puts it in magnet links, `Nip54InlineMetadata` in inline + * metadata, and `Nip47DeepLink` in the `callback`/`appname`/`value` parameters of NWC + * deep links. If Android encodes a title one way and iOS another, the two clients emit + * different bytes for the same event, and a wallet that round-trips a deep link built + * on one platform can fail on the other. + * + * The JVM/Android actual is `java.net.URLEncoder`/`URLDecoder` with UTF-8, so that is + * the reference every other target has to match. The expectations below are its + * `application/x-www-form-urlencoded` rules, which differ from plain RFC 3986 + * percent-encoding in exactly the three places a generic library gets "wrong": space, + * `*` and `~`. + */ +class UrlEncoderTest { + @Test + fun keepsTheUnreservedSet() { + // URLEncoder's dontNeedEncoding set is alphanumerics plus these four, and only + // these four. Note `*` survives and `~` does not — the opposite of RFC 3986. + assertEquals("abcXYZ019", UrlEncoder.encode("abcXYZ019")) + assertEquals("-_.*", UrlEncoder.encode("-_.*")) + } + + @Test + fun encodesSpaceAsPlus() { + // Form encoding, not %20. + assertEquals("hello+world", UrlEncoder.encode("hello world")) + assertEquals("a+b+c", UrlEncoder.encode("a b c")) + } + + @Test + fun encodesTildeAndTheOtherSubDelimiters() { + assertEquals("%7E", UrlEncoder.encode("~")) + assertEquals("%21", UrlEncoder.encode("!")) + assertEquals("%27", UrlEncoder.encode("'")) + assertEquals("%28%29", UrlEncoder.encode("()")) + assertEquals("%24%2C%3B", UrlEncoder.encode("$,;")) + } + + @Test + fun encodesUriPunctuationWithUppercaseHex() { + assertEquals("%3A%2F%2F", UrlEncoder.encode("://")) + assertEquals("%2B", UrlEncoder.encode("+")) + assertEquals("%3F%26%3D", UrlEncoder.encode("?&=")) + assertEquals("%40%23%25", UrlEncoder.encode("@#%")) + } + + @Test + fun encodesNonAsciiAsUtf8() { + assertEquals("%C3%A9", UrlEncoder.encode("é")) + assertEquals("caf%C3%A9", UrlEncoder.encode("café")) + assertEquals("%E2%82%AC", UrlEncoder.encode("€")) + // Outside the BMP: a surrogate pair has to encode as one 4-byte sequence. + assertEquals("%F0%9F%98%80", UrlEncoder.encode("😀")) + } + + @Test + fun decodesPlusAsSpace() { + assertEquals("hello world", UrlEncoder.decode("hello+world")) + assertEquals("hello world", UrlEncoder.decode("hello%20world")) + } + + @Test + fun decodesPercentEscapes() { + assertEquals("://", UrlEncoder.decode("%3A%2F%2F")) + assertEquals("~", UrlEncoder.decode("%7E")) + assertEquals("café", UrlEncoder.decode("caf%C3%A9")) + assertEquals("€", UrlEncoder.decode("%E2%82%AC")) + assertEquals("😀", UrlEncoder.decode("%F0%9F%98%80")) + // Lowercase hex decodes the same as uppercase. + assertEquals("é", UrlEncoder.decode("%c3%a9")) + } + + @Test + fun leavesUnescapedTextAlone() { + assertEquals("plain", UrlEncoder.decode("plain")) + assertEquals("-_.*~", UrlEncoder.decode("-_.*~")) + } + + @Test + fun roundTripsTheStringsThisIsActuallyUsedFor() { + // A torrent title (TorrentEvent) and a tracker URL. + val title = "Big Buck Bunny (2008) [1080p] ~ 60% done!" + assertEquals(title, UrlEncoder.decode(UrlEncoder.encode(title))) + + val tracker = "udp://tracker.example.org:1337/announce" + assertEquals("udp%3A%2F%2Ftracker.example.org%3A1337%2Fannounce", UrlEncoder.encode(tracker)) + assertEquals(tracker, UrlEncoder.decode(UrlEncoder.encode(tracker))) + + // An NWC pairing code (Nip47DeepLink.buildCallbackUri puts this in `value=`). + val pairing = + "nostr+walletconnect://b889ff5b1513b641e2a139f661a661364979c5beee91842f8f0ef42ab558e9d4" + + "?relay=wss%3A%2F%2Frelay.damus.io&secret=71a8c14c1407c113601079c4302dab36460f0ccd0ad506f1f2dc73b5100571c5" + assertEquals(pairing, UrlEncoder.decode(UrlEncoder.encode(pairing))) + + // A callback deep link (Nip47DeepLink.parseConnectUri reads this back). + val callback = "amethystnwc://callback" + assertEquals("amethystnwc%3A%2F%2Fcallback", UrlEncoder.encode(callback)) + assertEquals(callback, UrlEncoder.decode(UrlEncoder.encode(callback))) + } + + @Test + fun roundTripsEveryAsciiCharacter() { + val ascii = (0..127).map { it.toChar() }.joinToString("") + assertEquals(ascii, UrlEncoder.decode(UrlEncoder.encode(ascii))) + } +} diff --git a/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/UriParser.linux.kt b/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/UriParser.linux.kt index 448f4b8c8f..59c8a40b7d 100644 --- a/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/UriParser.linux.kt +++ b/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/UriParser.linux.kt @@ -122,9 +122,13 @@ actual class UriParser actual constructor( actual fun path(): String? = parsedPath /** - * Parsed once and reused, mirroring the JVM actual's lazy map. The previous version - * re-split the entire query string on every [getQueryParameter] call, so a URI read - * for four parameters was parsed four times. + * Parsed once and reused, mirroring the JVM actual's lazy map — the previous version + * re-split the entire query string on every [getQueryParameter] call. + * + * Decoded with [UrlEncoder.decode], which matches `URLDecoder.decode(.., "UTF-8")` — + * what the JVM actual calls. Skipping this is why NIP-47 failed on this target: + * `relay=wss%3A%2F%2Frelay.damus.io` reached `RelayUrlNormalizer` still encoded and + * came back "Invalid relay Url". */ private val queryParameters: Map> by lazy { parsedQuery?.ifBlank { null }?.let { query -> @@ -138,7 +142,7 @@ actual class UriParser actual constructor( } if (parts.size == 2) { - currentValue.add(percentDecode(parts[1])) + currentValue.add(UrlEncoder.decode(parts[1])) } else { currentValue.add("") } @@ -153,7 +157,7 @@ actual class UriParser actual constructor( keyValuePair.split('&').associate { paramValue -> val parts = paramValue.split("=", limit = 2) if (parts.size == 2) { - parts[0] to percentDecode(parts[1]) + parts[0] to UrlEncoder.decode(parts[1]) } else { parts[0] to "" // Handle parameters without a value } @@ -168,82 +172,3 @@ actual class UriParser actual constructor( actual fun fragments(): Map = parsedFragments } - -/** - * `java.net.URLDecoder.decode(value, "UTF-8")` — which is literally what the JVM actual - * calls — for a target with no `java.net`. - * - * This is the whole reason NIP-47 failed on linuxX64: the parser returned query values - * exactly as they appeared in the URI, so `relay=wss%3A%2F%2Frelay.damus.io` reached - * `RelayUrlNormalizer` still percent-encoded and came back "Invalid relay Url". Decoding - * belongs here rather than in each caller, because the JVM and Apple actuals both hand - * back decoded values and common code is written against that. - * - * Matches `URLDecoder` in the details that are observable: `+` becomes a space, a run of - * consecutive `%XX` is decoded as one UTF-8 sequence (so multi-byte characters survive), - * every other character passes through, and a malformed escape throws - * [IllegalArgumentException] rather than being silently kept — the same failure the JVM - * gives for the same input. - */ -private fun percentDecode(value: String): String { - // The short-circuit URLDecoder also makes: with nothing to change, return the - // original instance rather than rebuilding it. Most query values hit this. - if (value.indexOf('%') < 0 && value.indexOf('+') < 0) return value - - val result = StringBuilder(value.length) - var index = 0 - // Sized on first use for the longest run that could still follow, then reused — - // one allocation for the whole string, as in URLDecoder. - var escaped: ByteArray? = null - - while (index < value.length) { - when (val char = value[index]) { - '+' -> { - result.append(' ') - index++ - } - - '%' -> { - val buffer = escaped ?: ByteArray((value.length - index) / 3).also { escaped = it } - var count = 0 - while (index + 2 < value.length && value[index] == '%') { - buffer[count++] = decodeEscape(value, index) - index += 3 - } - if (index < value.length && value[index] == '%') { - throw IllegalArgumentException("URLDecoder: Incomplete trailing escape (%) pattern") - } - result.append(buffer.decodeToString(0, count)) - } - - else -> { - result.append(char) - index++ - } - } - } - - return result.toString() -} - -private fun decodeEscape( - value: String, - index: Int, -): Byte { - val high = hexDigit(value[index + 1]) - val low = hexDigit(value[index + 2]) - if (high < 0 || low < 0) { - throw IllegalArgumentException( - "URLDecoder: Illegal hex characters in escape (%) pattern - ${value.substring(index, index + 3)}", - ) - } - return ((high shl 4) or low).toByte() -} - -private fun hexDigit(char: Char): Int = - when (char) { - in '0'..'9' -> char - '0' - in 'a'..'f' -> char - 'a' + 10 - in 'A'..'F' -> char - 'A' + 10 - else -> -1 - } diff --git a/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/UrlEncoder.linux.kt b/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/UrlEncoder.linux.kt deleted file mode 100644 index 3b8e98f026..0000000000 --- a/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/UrlEncoder.linux.kt +++ /dev/null @@ -1,29 +0,0 @@ -/* - * Copyright (c) 2025 Vitor Pamplona - * - * Permission is hereby granted, free of charge, to any person obtaining a copy of - * this software and associated documentation files (the "Software"), to deal in - * the Software without restriction, including without limitation the rights to use, - * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the - * Software, and to permit persons to whom the Software is furnished to do so, - * subject to the following conditions: - * - * The above copyright notice and this permission notice shall be included in all - * copies or substantial portions of the Software. - * - * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR - * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS - * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR - * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN - * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION - * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. - */ -package com.vitorpamplona.quartz.utils - -import net.thauvin.erik.urlencoder.UrlEncoderUtil - -actual object UrlEncoder { - actual fun encode(value: String): String = UrlEncoderUtil.encode(value) - - actual fun decode(value: String): String = UrlEncoderUtil.decode(value) -} diff --git a/quartz/src/nativeMain/kotlin/com/vitorpamplona/quartz/utils/UrlEncoder.native.kt b/quartz/src/nativeMain/kotlin/com/vitorpamplona/quartz/utils/UrlEncoder.native.kt new file mode 100644 index 0000000000..2f1b1873e9 --- /dev/null +++ b/quartz/src/nativeMain/kotlin/com/vitorpamplona/quartz/utils/UrlEncoder.native.kt @@ -0,0 +1,190 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.quartz.utils + +/** + * Native actual for [UrlEncoder], shared by linuxX64 and every Apple target. + * + * ## Why this is not a library call any more + * + * Both native targets used to delegate to `net.thauvin.erik.urlencoder.UrlEncoderUtil`, + * which implements RFC 3986 percent-encoding. The JVM/Android actual is + * `java.net.URLEncoder`/`URLDecoder`, which implements + * `application/x-www-form-urlencoded`. Those are different specifications, and the + * difference was observable in three places: + * + * ``` + * JVM/Android UrlEncoderUtil + * encode(" ") "+" "%20" + * encode("*") "*" "%2A" + * decode("a+b") "a b" "a+b" + * ``` + * + * That is not cosmetic. [encode] builds strings that leave the device — `TorrentEvent` + * puts it in magnet links, `Nip54InlineMetadata` in inline metadata, `Nip47DeepLink` in + * the `callback`, `appname` and `value` parameters of NWC deep links — so Android and + * iOS were emitting different bytes for the same title. The decode row is worse than + * cosmetic: a magnet link or deep link written by Android carries `+` for its spaces, + * and reading it on iOS produced a string with literal plus signs instead of spaces, no + * error anywhere. + * + * So this matches `URLEncoder`/`URLDecoder` exactly instead: the unreserved set is + * alphanumerics plus `-`, `_`, `.` and `*` (note `*` survives and `~` does not — the + * opposite of RFC 3986), space encodes to `+`, everything else to uppercase `%XX` of + * its UTF-8 bytes, and decoding maps `+` back to a space. `UrlEncoderTest` in + * `commonTest` pins all of it against the JVM on every target. + * + * Both directions short-circuit the way the `java.net` pair does: a string with nothing + * to change is returned as-is rather than rebuilt. + * + * One deliberate edge difference: an *unpaired* UTF-16 surrogate encodes as `%EF%BF%BD` + * (Kotlin's replacement character) where the JVM gives `%3F`. Nostr content is + * well-formed UTF-16, and chasing it would cost a scan on every call. + */ +actual object UrlEncoder { + private const val HEX = "0123456789ABCDEF" + + actual fun encode(value: String): String { + var index = 0 + while (index < value.length && isUnreserved(value[index])) index++ + if (index == value.length) return value + + val result = StringBuilder(value.length + ESCAPE_HEADROOM) + result.append(value, 0, index) + + while (index < value.length) { + val char = value[index] + when { + isUnreserved(char) -> { + result.append(char) + index++ + } + + char == ' ' -> { + result.append('+') + index++ + } + + else -> { + // Escaped as a run rather than character by character, so a surrogate + // pair becomes one 4-byte sequence instead of two malformed 3-byte ones. + val start = index + do { + index++ + } while (index < value.length && !isUnreserved(value[index]) && value[index] != ' ') + appendEscaped(result, value, start, index) + } + } + } + + return result.toString() + } + + actual fun decode(value: String): String { + if (value.indexOf('%') < 0 && value.indexOf('+') < 0) return value + + val result = StringBuilder(value.length) + var index = 0 + // Sized on first use for the longest run that could still follow, then reused — + // one allocation for the whole string, as in URLDecoder. + var escaped: ByteArray? = null + + while (index < value.length) { + when (val char = value[index]) { + '+' -> { + result.append(' ') + index++ + } + + '%' -> { + val buffer = escaped ?: ByteArray((value.length - index) / 3).also { escaped = it } + var count = 0 + while (index + 2 < value.length && value[index] == '%') { + buffer[count++] = decodeEscape(value, index) + index += 3 + } + if (index < value.length && value[index] == '%') { + throw IllegalArgumentException("URLDecoder: Incomplete trailing escape (%) pattern") + } + // Decoded as a run so a multi-byte UTF-8 sequence survives. + result.append(buffer.decodeToString(0, count)) + } + + else -> { + result.append(char) + index++ + } + } + } + + return result.toString() + } + + /** `URLEncoder`'s `dontNeedEncoding` set: alphanumerics plus these four, and only these. */ + private fun isUnreserved(char: Char): Boolean = + char in 'a'..'z' || + char in 'A'..'Z' || + char in '0'..'9' || + char == '-' || + char == '_' || + char == '.' || + char == '*' + + private fun appendEscaped( + result: StringBuilder, + value: String, + start: Int, + end: Int, + ) { + val bytes = value.substring(start, end).encodeToByteArray() + for (byte in bytes) { + val code = byte.toInt() + result.append('%') + result.append(HEX[(code shr 4) and 0xF]) + result.append(HEX[code and 0xF]) + } + } + + private fun decodeEscape( + value: String, + index: Int, + ): Byte { + val high = hexDigit(value[index + 1]) + val low = hexDigit(value[index + 2]) + if (high < 0 || low < 0) { + throw IllegalArgumentException( + "URLDecoder: Illegal hex characters in escape (%) pattern - ${value.substring(index, index + 3)}", + ) + } + return ((high shl 4) or low).toByte() + } + + private fun hexDigit(char: Char): Int = + when (char) { + in '0'..'9' -> char - '0' + in 'a'..'f' -> char - 'a' + 10 + in 'A'..'F' -> char - 'A' + 10 + else -> -1 + } + + /** Enough for a handful of escapes before the builder has to grow. */ + private const val ESCAPE_HEADROOM = 16 +} From 9e2859f719e038150347d274cc56aa0a08da97e8 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 1 Sep 2026 21:01:16 +0000 Subject: [PATCH 6/6] fix(quartz): stripe from the bucket, not from unrelated hash bits MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Audit finding, and a real defect in the striped table two commits back. A striped hash table is only sound when the stripe is a function of the bucket. lockFor picked bits 16-19 of the hash while the bucket index used the low bits, so the 16 locks did not partition the table: two keys could share a bucket while holding different locks, and two writers would then read the same chain head and both publish over it. One insert silently disappears while entryCount counts both. The same window loses an overwrite, and loses entries through remove's chain rebuild. That is precisely the class of bug this work set out to remove from the copy-on-write version it replaced. Stripe now comes from `hash and (STRIPES - 1)`. Because STRIPES and every capacity are powers of two with STRIPES <= capacity, those are exactly the low bits of the bucket index, so same bucket implies same stripe at every size. It stays derived from the hash rather than the capacity, so a key keeps its stripe across a resize, which is what lets growTable exclude writers by taking all of them. INITIAL_CAPACITY is now defined as STRIPES so raising one cannot silently break the invariant. That definition also fixes a memory regression the audit caught: the table allocated 1024 slots eagerly, about 8 KB, per instance. LargeCache is not only the one big LocalCache — EphemeralRoom, RelaySession, PoolRequests and others build one per room, per connection and per subscription set, so a client holds hundreds that stay nearly empty. An empty instance goes from ~8 KB to ~970 bytes. Growth is geometric, so a table that does fill to 100k pays the same ~2n node rebuilds either way; re-measuring the shipped code confirms it (fill 16ms, overwrite 4ms, reads 1ms, 20 scans 16ms, mixed 86ms, 1 GC — unchanged within noise). The KDoc table is updated to those numbers. Adds LargeCacheStripingTest, which builds keys that share a bucket while differing in bits 16-19 and drives four workers at them behind a start barrier, with few enough buckets that chains grow long and each insert holds its lock for a while. It is documented for what it is: a stress test of the concurrent same-bucket path, not a deterministic reproducer — it did not fail against the broken striping in the runs attempted, which makes that race rare rather than absent. The fix rests on reading the stripe selection against the bucket index, not on a red test. Remaining known cost, noted in the KDoc rather than changed here: those ~970 bytes are nearly all the 16 PlatformLocks, two objects each. Folding them into one AtomicIntArray would reach ~250 bytes, but hand-rolling the spin wants its own review rather than a change on the way to merge. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01HxQ1QuyzSkR38iFHbREjoS --- .../quartz/utils/cache/StripedHashMap.kt | 49 ++++- .../utils/cache/LargeCacheStripingTest.kt | 183 ++++++++++++++++++ 2 files changed, 224 insertions(+), 8 deletions(-) create mode 100644 quartz/src/linuxTest/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCacheStripingTest.kt diff --git a/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/StripedHashMap.kt b/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/StripedHashMap.kt index ec9caba654..b50976cafa 100644 --- a/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/StripedHashMap.kt +++ b/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/cache/StripedHashMap.kt @@ -58,14 +58,23 @@ import kotlin.concurrent.atomics.ExperimentalAtomicApi * copy-on-write* n/a n/a n/a n/a n/a n/a n/a * HAMT + CAS 70 78 6 117 676 24 67MB * lock + HashMap 13 6 3 71 1197 36 51MB - * this 15 3 3 13 64 1 43MB + * this 16 4 1 16 86 1 41MB * ``` * * (*copy-on-write is off the scale: 20k entries alone took 18s to fill.) "mixed" is a * full fill with a whole-table scan every 1000 writes, which is the shape `LocalCache` * actually has — arriving events interleaved with feeds filtering the whole cache. - * Figures are one representative run of several; they were stable to within ~10%, - * except the HAMT's scan column, which wandered between 115ms and 190ms. + * Figures are one representative run of several; they were stable to within ~10%, except + * the HAMT's scan column, which wandered between 115ms and 190ms. This row was + * re-measured on the shipped code, after the stripe-selection fix and after + * [INITIAL_CAPACITY] dropped to [STRIPES] — neither moved it out of the noise. + * + * An *empty* instance costs ~970 bytes, nearly all of it the 16 [PlatformLock]s (two + * objects each). That is well above the ~50 bytes an empty map used to cost, and it is + * charged to every one of the hundreds of small caches a client holds. Folding the stripe + * locks into a single `AtomicIntArray` would take it to ~250 bytes and is the obvious next + * step if it ever shows up in a heap profile; it is left alone here because hand-rolling + * the spin is exactly the kind of change that wants its own review. * * ## Concurrency contract * @@ -111,10 +120,23 @@ internal class StripedHashMap { } /** - * Derived from the hash alone, never from the table size, so a key keeps the same - * stripe across a resize. + * **The stripe must be a function of the bucket**, or the locks do not partition the + * table and the whole design is unsound: two keys could share a bucket while holding + * different locks, so two writers would read the same chain head and both publish + * over it, silently dropping one insert. + * + * It is a function of the bucket here because [STRIPES] and every table capacity are + * powers of two with `STRIPES <= capacity`, so `hash and (STRIPES - 1)` is exactly the + * low bits of `hash and (capacity - 1)`. Same bucket therefore implies same stripe, at + * every size. Using the hash and not the capacity also keeps a key on one stripe + * across a resize, which is what lets [growTable] exclude writers by taking all of + * them. + * + * An earlier version took bits 16-19 instead. That is still resize-stable and still + * spreads well, which is why it looked right — but it is not derived from the bucket, + * so it broke the invariant above. */ - private fun lockFor(hash: Int) = locks[(hash ushr 16) and (STRIPES - 1)] + private fun lockFor(hash: Int) = locks[hash and (STRIPES - 1)] fun size(): Int = entryCount.load() @@ -301,8 +323,19 @@ internal class StripedHashMap { */ private const val STRIPES = 16 - /** Sized to carry a warm cache's first few thousand entries without a resize. */ - private const val INITIAL_CAPACITY = 1024 + /** + * Tied to [STRIPES] rather than chosen independently: the stripe is only a + * function of the bucket while `STRIPES <= capacity` (see [lockFor]), and defining + * it this way makes that impossible to break by raising [STRIPES] alone. + * + * Kept small on purpose. `LargeCache` is not only the one big `LocalCache` + * instance: `EphemeralRoom`, `RelaySession`, `PoolRequests` and friends each build + * one per room, per connection and per subscription set, so a client holds + * hundreds of them and most stay nearly empty. Starting at 1024 slots charged every + * one of those ~8 KB it would never use. Growth is geometric, so a table that does + * fill to 100k pays the same ~2n node rebuilds in total either way. + */ + private const val INITIAL_CAPACITY = STRIPES private const val MAX_CAPACITY = 1 shl 30 } diff --git a/quartz/src/linuxTest/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCacheStripingTest.kt b/quartz/src/linuxTest/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCacheStripingTest.kt new file mode 100644 index 0000000000..ec2b1901d4 --- /dev/null +++ b/quartz/src/linuxTest/kotlin/com/vitorpamplona/quartz/utils/cache/LargeCacheStripingTest.kt @@ -0,0 +1,183 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.quartz.utils.cache + +import kotlin.concurrent.atomics.AtomicInt +import kotlin.concurrent.atomics.ExperimentalAtomicApi +import kotlin.native.concurrent.ObsoleteWorkersApi +import kotlin.native.concurrent.TransferMode +import kotlin.native.concurrent.Worker +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertNotNull + +/** + * Stresses the one path the rest of the suite never reaches: **many writers inserting into + * the same bucket at the same time.** + * + * A striped table is only sound if the stripe is a function of the bucket. Derive the two + * from different parts of the hash and the locks stop partitioning the table: two keys can + * share a bucket while holding different locks, so two writers read the same chain head and + * both publish over it, and one insert vanishes. The first cut of [StripedHashMap] did + * exactly that — bucket from the low bits, stripe from bits 16-19 — and nothing here caught + * it, because every other suite uses small `Int` keys or a constant `hashCode`, which all + * collapse onto stripe 0 and serialise by accident. + * + * So the keys are built on purpose: groups of four sharing their low 16 bits, which puts a + * group in one bucket at any table size this reaches, while differing in bits 16-19 — the + * part the broken version striped on. Only [BUCKETS] buckets are used, so chains run to + * hundreds of nodes and each insert holds its lock for a while, which is what widens the + * window; with well-spread keys it is a few nanoseconds and nothing is observable. + * + * Being straight about what this is: a stress test, not a deterministic reproducer. It did + * not fail against the broken striping in the runs attempted, which says that race is rare + * rather than absent — the defect is a plain lost update, provable by reading + * [StripedHashMap]'s stripe selection against its bucket index, and was fixed on that + * basis. What this test is worth is being the only coverage of concurrent same-bucket + * inserts at all, and it would catch a coarser regression. + */ +@OptIn(ObsoleteWorkersApi::class, ExperimentalAtomicApi::class) +class LargeCacheStripingTest { + /** + * `hashOf` in [StripedHashMap] spreads with `h xor (h ushr 16)`, so to land on a final + * hash of `(member shl 16) or bucket` the raw hashCode has to be + * `(member shl 16) or (bucket xor member)`. The low 16 bits are then exactly `bucket`, + * shared by all four members of a group at every table size this test reaches. + * + * Only [BUCKETS] distinct buckets are used, so chains run to hundreds of nodes. That + * matters: a writer walks its chain *holding the stripe lock*, so a long chain is what + * makes the window wide enough for two writers on two different locks to overlap + * inside the same bucket. With well-spread keys the window is a few nanoseconds and + * the defect hides. + */ + private data class GroupedKey( + val group: Int, + val member: Int, + ) { + override fun hashCode(): Int { + val bucket = group and (BUCKETS - 1) + return (member shl 16) or (bucket xor member) + } + } + + /** Runs [job] on [WORKERS] threads released together, so they actually contend. */ + private fun inParallel(job: (workerId: Int) -> Unit) { + val ready = AtomicInt(0) + val go = AtomicInt(0) + val workers = List(WORKERS) { Worker.start() } + + val futures = + workers.mapIndexed { id, worker -> + worker.execute(TransferMode.SAFE, { Triple(job, id, ready to go) }) { (block, workerId, gates) -> + val (readyGate, goGate) = gates + readyGate.fetchAndAdd(1) + while (goGate.load() == 0) { } + block(workerId) + } + } + + while (ready.load() < WORKERS) { } + go.store(1) + + futures.forEach { it.result } + workers.forEach { it.requestTermination().result } + } + + @Test + fun concurrentInsertsAcrossSharedBucketsKeepEveryEntry() { + repeat(ROUNDS) { round -> insertRound(round) } + } + + private fun insertRound(round: Int) { + val cache = LargeCache() + + inParallel { worker -> + for (group in 0 until GROUPS) { + cache.put(GroupedKey(group, worker), group) + } + } + + val total = GROUPS * WORKERS + + val seen = mutableSetOf() + cache.forEach { key, _ -> seen.add(key) } + + assertEquals(total, seen.size, "round $round: iteration lost or duplicated entries in a shared bucket") + assertEquals(total, cache.size(), "round $round: size() disagrees with what the table holds") + + for (group in 0 until GROUPS) { + for (member in 0 until WORKERS) { + assertEquals( + group, + assertNotNull( + cache.get(GroupedKey(group, member)), + "round $round: entry ($group, $member) was dropped by a concurrent insert", + ), + ) + } + } + } + + @Test + fun concurrentCreateIfAbsentAcrossSharedBucketsReportsOneInsertEach() { + repeat(ROUNDS) { createIfAbsentRound() } + } + + private fun createIfAbsentRound() { + val cache = LargeCache() + val contended = GROUPS / 4 + + // Every worker races for the same keys this time, so a lost update shows up as a + // duplicate in the chain rather than a missing entry. + inParallel { _ -> + for (group in 0 until contended) { + for (member in 0 until WORKERS) { + cache.createIfAbsent(GroupedKey(group, member)) { group } + } + } + } + + val total = contended * WORKERS + val seen = mutableSetOf() + var visited = 0 + cache.forEach { key, _ -> + seen.add(key) + visited++ + } + + assertEquals(total, visited, "a key was inserted twice into the same chain") + assertEquals(total, seen.size) + assertEquals(total, cache.size()) + } + + companion object { + private const val WORKERS = 4 + + /** Large enough that the four workers overlap for essentially the whole run. */ + private const val GROUPS = 4_096 + + /** Few enough that chains grow long and every insert holds its lock for a while. */ + private const val BUCKETS = 64 + + /** Repeated on a fresh table, because only the *insert* path can lose a write. */ + private const val ROUNDS = 20 + } +}