fix(quartz): back Apple LargeCache with StripedHashMap, not CacheMap

NostrSignerRemotePrivateZapTest failed on iosSimulatorArm64 with
kotlin.ConcurrentModificationException. Apple's LargeCache wrapped
charlietap's CacheMap, a left-right map whose entries/keys/values return
live views of an inner HashMap and drop the read guard before the caller
iterates. Every bulk op (filter, map, keys(), even forEach's
entries.toList() "snapshot") therefore walked a HashMap a writer could be
mutating, and Kotlin/Native's HashMap throws CME when that happens. The
bunker round-trip has NostrClient's pool scanning its LargeCaches while
subscribe/publish mutate them from other threads.

Reproduced on linuxX64 (same Kotlin/Native HashMap) by racing CacheMap
scans against put/remove workers: kotlin.ConcurrentModificationException.

Linux already had a purpose-built replacement, StripedHashMap: lock-free,
weakly consistent reads and striped-lock writes, so scans never throw and
never copy. It only needs PlatformLock, which has a parking
(NSRecursiveLock) Apple actual. Move it and the LargeCache /
ConcurrentHashCache actuals from linuxMain to nativeMain so Apple uses
them too, move their tests to nativeTest so they run on iOS, and drop the
now-unused cachemap dependency.

Adds LargeCacheConcurrencyTest.scansNeverThrowWhileKeysComeAndGo, which
covers the filter/map/keys/values/forEach scans against concurrent puts
and removes.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01PcC7xe3bTb9fnJhESGukuD
This commit is contained in:
Claude
2026-10-02 15:43:13 +00:00
parent 6a699dcc4c
commit 50486341a6
15 changed files with 65 additions and 497 deletions
@@ -70,7 +70,7 @@ class ChessEventCollector(
val startEvent: StateFlow<JesterEvent?> = _startEvent.asStateFlow()
// Move events (deduplicated by event ID). String keys are Comparable, so
// LargeCache (ConcurrentSkipListMap on JVM, CacheMap on Apple) works.
// LargeCache (ConcurrentSkipListMap on JVM, StripedHashMap on native) works.
private val moves = LargeCache<String, JesterEvent>()
// Track all processed event IDs for fast deduplication. A plain HashSet
-2
View File
@@ -10,7 +10,6 @@ accompanistAdaptive = "0.37.3"
# LargeCache/ConcurrentHashCache — they read those views directly and through the
# stdlib Map operators (filter/map/groupBy/associate/count). Bumping needs a port to
# the new forEach/forEachKey/forEachValue API, verified on an Apple target.
cachemapVersion = "0.2.4"
composeMultiplatform = "1.12.1"
activityCompose = "1.13.0"
# 9.4.0's PerModuleBundleTask rejects an AAB entry whose name contains a colon, which every
@@ -174,7 +173,6 @@ androidx-ui-test-junit4 = { group = "androidx.compose.ui", name = "ui-test-junit
androidx-ui-tooling = { group = "androidx.compose.ui", name = "ui-tooling" }
androidx-ui-tooling-preview = { group = "androidx.compose.ui", name = "ui-tooling-preview" }
audiowaveform = { group = "com.github.lincollincol", name = "compose-audiowaveform", version.ref = "audiowaveform" }
charlietap-cachemap = { module = "io.github.charlietap:cachemap", version.ref = "cachemapVersion" }
coil-compose = { group = "io.coil-kt.coil3", name = "coil-compose", version.ref = "coil" }
coil-gif = { group = "io.coil-kt.coil3", name = "coil-gif", version.ref = "coil" }
coil-svg = { group = "io.coil-kt.coil3", name = "coil-svg", version.ref = "coil" }
-1
View File
@@ -296,7 +296,6 @@ kotlin {
create("appleMain") {
dependsOn(nativeMain)
dependencies {
implementation(libs.charlietap.cachemap)
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")
@@ -1,40 +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.cache
import io.github.charlietap.cachemap.cacheMapOf
actual class ConcurrentHashCache<K : Any, V : Any> {
private val map = cacheMapOf<K, V>()
actual fun get(key: K): V? = map[key]
actual fun put(
key: K,
value: V,
) {
map.put(key, value)
}
actual fun size(): Int = map.size
actual fun clear() = map.clear()
}
@@ -1,432 +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.cache
import io.github.charlietap.cachemap.CacheMap
import io.github.charlietap.cachemap.cacheMapOf
// An implementation of a Threadsafe map, using CacheMap.
// Investigating a Swift-based alternative(for now)
actual class LargeCache<K, V> : ICacheOperations<K, V> {
private val concurrentMap = cacheMapOf<K, V>()
actual fun keys(): Set<K> = concurrentMap.keys
actual fun values(): Iterable<V> = concurrentMap.values
actual fun get(key: K): V? = concurrentMap[key]
actual fun remove(key: K): V? = concurrentMap.remove(key)
actual fun isEmpty(): Boolean = concurrentMap.isEmpty()
actual fun clear() {
concurrentMap.clear()
}
actual fun containsKey(key: K): Boolean = concurrentMap.containsKey(key)
actual fun put(
key: K,
value: V,
) {
concurrentMap.put(key, value)
}
actual fun getOrCreate(
key: K,
builder: (key: K) -> V,
): V {
val value = concurrentMap.get(key)
return if (value != null) {
value
} else {
val newObject = builder(key)
concurrentMap.put(key, newObject)
concurrentMap[key] ?: newObject
}
}
actual fun createIfAbsent(
key: K,
builder: (key: K) -> V,
): Boolean {
val value = concurrentMap.get(key)
return if (value != null) {
false
} else {
val newObject = builder(key)
concurrentMap.put(key, newObject)
concurrentMap[key] != null
}
}
actual override fun size(): Int = concurrentMap.size
actual override fun forEach(consumer: ICacheBiConsumer<K, V>) {
// Take a snapshot of entries to avoid ConcurrentModificationException
// when the map is modified during iteration (e.g., NostrClient.syncFilters
// iterating while subscriptions are added from another coroutine).
concurrentMap.entries.toList().forEach { consumer.accept(it.key, it.value) }
}
actual override fun filter(consumer: CacheCollectors.BiFilter<K, V>): List<V> =
concurrentMap
.filter { consumer.filter(it.key, it.value) }
.values
.toList()
actual override fun filterIntoSet(consumer: CacheCollectors.BiFilter<K, V>): Set<V> =
concurrentMap
.filter { consumer.filter(it.key, it.value) }
.values
.toSet()
actual override fun <R> map(consumer: CacheCollectors.BiNotNullMapper<K, V, R>): List<R> = concurrentMap.map { consumer.map(it.key, it.value) }
actual override fun <R> mapNotNull(consumer: CacheCollectors.BiMapper<K, V, R?>): List<R> = concurrentMap.mapNotNull { consumer.map(it.key, it.value) }
actual override fun <R> mapNotNullIntoSet(consumer: CacheCollectors.BiMapper<K, V, R?>): Set<R> = mapNotNull(consumer).toSet()
actual override fun <R> mapFlatten(consumer: CacheCollectors.BiMapper<K, V, Collection<R>?>): List<R> = concurrentMap.flatMap { entry -> consumer.map(entry.key, entry.value) ?: emptyList() }
actual override fun <R> mapFlattenIntoSet(consumer: CacheCollectors.BiMapper<K, V, Collection<R>?>): Set<R> = mapFlatten(consumer).toSet()
actual override fun maxOrNullOf(
filter: CacheCollectors.BiFilter<K, V>,
comparator: Comparator<V>,
): V? {
// return concurrentMap.maxOfWithOrNull(
// comparator,
// selector = {
// if (filter.filter(it.key, it.value)) it.value else concurrentMap.getValue(it.key)
// }
// )
return concurrentMap.maxOrNullOf(filter, comparator)
}
actual override fun sumOf(consumer: CacheCollectors.BiSumOf<K, V>): Int {
return concurrentMap.sumOf(consumer)
// return concurrentMap.map { consumer.map(it.key, it.value) }.sum()
}
actual override fun sumOfLong(consumer: CacheCollectors.BiSumOfLong<K, V>): Long = concurrentMap.sumOfLong(consumer)
actual override fun <R> groupBy(consumer: CacheCollectors.BiNotNullMapper<K, V, R>): Map<R, List<V>> = concurrentMap.groupBy(consumer)
actual override fun <R> countByGroup(consumer: CacheCollectors.BiNotNullMapper<K, V, R>): Map<R, Int> = concurrentMap.countByGroup(consumer)
actual override fun <R> sumByGroup(
groupMap: CacheCollectors.BiNotNullMapper<K, V, R>,
sumOf: CacheCollectors.BiNotNullMapper<K, V, Long>,
): Map<R, Long> = concurrentMap.sumByGroup(groupMap, sumOf)
actual override fun count(consumer: CacheCollectors.BiFilter<K, V>): Int = concurrentMap.count { consumer.filter(it.key, it.value) }
actual override fun <T, U> associate(transform: (K, V) -> Pair<T, U>): Map<T, U> = concurrentMap.associate(transform)
actual override fun <U> associateWith(transform: (K, V) -> U?): Map<K, U?> = concurrentMap.associateWith(transform)
actual override fun filter(
from: K,
to: K,
consumer: CacheCollectors.BiFilter<K, V>,
): List<V> {
val transientList = concurrentMap.subMapAlt(from, to)
return transientList.filter { consumer.filter(it.key, it.value) }.values.toList()
}
actual override fun filterIntoSet(
from: K,
to: K,
consumer: CacheCollectors.BiFilter<K, V>,
): Set<V> = filter(from, to, consumer).toSet()
actual override fun <R> map(
from: K,
to: K,
consumer: CacheCollectors.BiNotNullMapper<K, V, R>,
): List<R> {
val transientList = concurrentMap.subMapAlt(from, to)
return transientList.map { consumer.map(it.key, it.value) }
}
actual override fun <R> mapNotNull(
from: K,
to: K,
consumer: CacheCollectors.BiMapper<K, V, R?>,
): List<R> = concurrentMap.subMapAlt(from, to).mapNotNull { consumer.map(it.key, it.value) }
actual override fun <R> mapNotNullIntoSet(
from: K,
to: K,
consumer: CacheCollectors.BiMapper<K, V, R?>,
): Set<R> = mapNotNull(from, to, consumer).toSet()
actual override fun <R> mapFlatten(
from: K,
to: K,
consumer: CacheCollectors.BiMapper<K, V, Collection<R>?>,
): List<R> = concurrentMap.subMapAlt(from, to).flatMap { consumer.map(it.key, it.value) as Iterable<R> }
actual override fun <R> mapFlattenIntoSet(
from: K,
to: K,
consumer: CacheCollectors.BiMapper<K, V, Collection<R>?>,
): Set<R> = mapFlatten(from, to, consumer).toSet()
actual override fun maxOrNullOf(
from: K,
to: K,
filter: CacheCollectors.BiFilter<K, V>,
comparator: Comparator<V>,
): V? {
val transient = concurrentMap.subMapAlt(from, to)
return transient.maxOrNullOf(filter, comparator)
}
actual override fun sumOf(
from: K,
to: K,
consumer: CacheCollectors.BiSumOf<K, V>,
): Int = concurrentMap.subMapAlt(from, to).sumOf(consumer)
actual override fun sumOfLong(
from: K,
to: K,
consumer: CacheCollectors.BiSumOfLong<K, V>,
): Long = concurrentMap.subMapAlt(from, to).sumOfLong(consumer)
actual override fun <R> groupBy(
from: K,
to: K,
consumer: CacheCollectors.BiNotNullMapper<K, V, R>,
): Map<R, List<V>> = concurrentMap.subMapAlt(from, to).groupBy(consumer)
actual override fun <R> countByGroup(
from: K,
to: K,
consumer: CacheCollectors.BiNotNullMapper<K, V, R>,
): Map<R, Int> = concurrentMap.subMapAlt(from, to).countByGroup(consumer)
actual override fun <R> sumByGroup(
from: K,
to: K,
groupMap: CacheCollectors.BiNotNullMapper<K, V, R>,
sumOf: CacheCollectors.BiNotNullMapper<K, V, Long>,
): Map<R, Long> = concurrentMap.subMapAlt(from, to).sumByGroup(groupMap, sumOf)
actual override fun count(
from: K,
to: K,
consumer: CacheCollectors.BiFilter<K, V>,
): Int = concurrentMap.subMapAlt(from, to).count { consumer.filter(it.key, it.value) }
actual override fun <T, U> associate(
from: K,
to: K,
transform: (K, V) -> Pair<T, U>,
): Map<T, U> = concurrentMap.subMapAlt(from, to).associate(transform)
actual override fun <U> associateWith(
from: K,
to: K,
transform: (K, V) -> U?,
): Map<K, U?> = concurrentMap.subMapAlt(from, to).associateWith(transform)
actual override fun joinToString(
separator: CharSequence,
prefix: CharSequence,
postfix: CharSequence,
limit: Int,
truncated: CharSequence,
transform: ((K, V) -> CharSequence)?,
): String {
val buffer = StringBuilder()
buffer.append(prefix)
var count = 0
forEach { key, value ->
val str = if (transform != null) transform(key, value) else ""
if (str.isNotEmpty()) {
if (++count > 1) buffer.append(separator)
if (limit < 0 || count <= limit) {
when {
transform != null -> buffer.append(str)
else -> buffer.append("$key $value")
}
} else {
return@forEach
}
}
}
if (limit >= 0 && count > limit) buffer.append(truncated)
buffer.append(postfix)
return buffer.toString()
}
}
// Different subMap implementations below. Investigating their performance for now.
fun <K, V> CacheMap<K, V>.subMapSlow(
from: K,
to: K,
toInclusive: Boolean = true,
): Map<K, V> {
val transientList = toList()
val transientSubList =
transientList.subList(
fromIndex = transientList.indexOf(Pair(from, getValue(from))),
toIndex = transientList.indexOf(Pair(to, getValue(to))),
)
val completeSubList = transientSubList + Pair(to, getValue(to))
return if (toInclusive) completeSubList.toMap() else transientSubList.toMap()
}
fun <K, V> CacheMap<K, V>.subMapAlt(
from: K,
to: K,
toInclusive: Boolean = true,
): Map<K, V> {
val resultMap = hashMapOf<K, V>()
val keySet = keys
val fromIndex = keySet.indexOf(from)
val toIndex = keySet.indexOf(to)
for (index in fromIndex until toIndex) {
val correspondingEntry = entries.elementAt(index)
resultMap[correspondingEntry.key] = correspondingEntry.value
}
if (toInclusive) {
val correspondingToEntry = entries.elementAt(toIndex)
resultMap[correspondingToEntry.key] = correspondingToEntry.value
}
return resultMap
}
/**
* The following functions below are (re)implementations for the ICacheOperations
* interface. A lot of it is copying and pasting, with modifications to make it work
* consistently.
*/
fun <K, V> Map<K, V>.maxOrNullOf(
filter: CacheCollectors.BiFilter<K, V>,
comparator: Comparator<V>,
): V? {
var maxK: K? = null
var maxV: V? = null
forEach {
if (filter.filter(it.key, it.value)) {
if (maxK == null || (maxV != null && comparator.compare(it.value, maxV) > 0)) {
maxK = it.key
maxV = it.value
}
}
}
val finalMaxK: K? = maxK
val finalMaxV: V? = maxV
return finalMaxV
}
fun <K, V> Map<K, V>.sumOf(consumer: CacheCollectors.BiSumOf<K, V>): Int {
var sum = 0
forEach { sum += consumer.map(it.key, it.value) }
return sum
}
fun <K, V> Map<K, V>.sumOfLong(consumer: CacheCollectors.BiSumOfLong<K, V>): Long {
var sum = 0L
forEach { sum += consumer.map(it.key, it.value) }
return sum
}
fun <K, V, R> Map<K, V>.groupBy(consumer: CacheCollectors.BiNotNullMapper<K, V, R>): Map<R, List<V>> {
val results = HashMap<R, ArrayList<V>>()
forEach {
val group = consumer.map(it.key, it.value)
val list = results[group]
if (list == null) {
val answer = ArrayList<V>()
answer.add(it.value)
results[group] = answer
} else {
list.add(it.value)
}
}
return results
}
fun <K, V, R> Map<K, V>.countByGroup(consumer: CacheCollectors.BiNotNullMapper<K, V, R>): Map<R, Int> {
val results = HashMap<R, Int>()
forEach {
val group = consumer.map(it.key, it.value)
val count = results[group]
if (count == null) {
results[group] = 1
} else {
results[group] = count + 1
}
}
return results
}
fun <K, V, R> Map<K, V>.sumByGroup(
groupMap: CacheCollectors.BiNotNullMapper<K, V, R>,
sumOf: CacheCollectors.BiNotNullMapper<K, V, Long>,
): Map<R, Long> {
val results = HashMap<R, Long>()
forEach {
val group = groupMap.map(it.key, it.value)
val sum = results[group]
if (sum == null) {
results[group] = sumOf.map(it.key, it.value)
} else {
results[group] = sum + sumOf.map(it.key, it.value)
}
}
return results
}
fun <K, V, T, U> Map<K, V>.associate(transform: (K, V) -> Pair<T, U>): Map<T, U> {
val results: LinkedHashMap<T, U> = LinkedHashMap(size)
forEach {
val pair = transform(it.key, it.value)
results[pair.first] = pair.second
}
return results
}
fun <K, V, U> Map<K, V>.associateWith(transform: (K, V) -> U?): Map<K, U?> {
val results: LinkedHashMap<K, U?> = LinkedHashMap(size)
forEach {
results[it.key] = transform(it.key, it.value)
}
return results
}
@@ -106,7 +106,7 @@ class RelayAuthenticator(
) : IAuthStatus {
// Connection callbacks fire on the per-relay OkHttp dispatcher thread, so
// this state is mutated concurrently — LargeCache wraps a platform-tuned
// concurrent map (ConcurrentSkipListMap on jvmAndroid, CacheMap on Apple).
// concurrent map (ConcurrentSkipListMap on jvmAndroid, StripedHashMap on native).
//
// This stays mutable because RelayAuthStatus carries an LruCache that has
// to be addressable from the dispatcher thread. The Compose-observable
@@ -30,8 +30,8 @@ package com.vitorpamplona.quartz.utils.cache
* duplicate check in
* [com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.CachingEventDecoder].
*
* Actuals: JVM/Android → `ConcurrentHashMap`; Apple → `CacheMap` (same
* backing as [LargeCache]); Linux → copy-on-write (CI-only target).
* Actuals: JVM/Android → `ConcurrentHashMap`; Apple and Linux → a striped-lock
* chained hash table (`StripedHashMap`, same backing as [LargeCache]).
*/
expect class ConcurrentHashCache<K : Any, V : Any>() {
fun get(key: K): V?
@@ -62,20 +62,18 @@ import kotlin.time.TimeSource
* two HashMaps per put, 60k live ratio 2.28 - 5.21
* ```
*
* The last row is what Apple targets actually run: LargeCache there wraps
* The last row is what Apple targets used to run: LargeCache there wrapped
* charlietap's CacheMap, whose LeftRight `mutate` applies each write to both
* of its two maps under a lock — O(1), no copying. It still crossed the 5.0
* threshold on a loaded machine, which is how this failed on the iOS simulator
* without anything being wrong with the outbox. The `retains nothing` row is
* the control: same allocations, nothing kept alive, ratio flat.
*
* So on Apple the O(1) guarantee comes from the data structure by
* construction and does not need pinning here; on JVM/Android it comes from
* LargeCache being a ConcurrentHashMap, and a regression in this class
* (someone reintroducing a copy-on-write map, or a per-publish full scan)
* shows up cleanly. Note that Kotlin/Native's linuxX64 LargeCache IS
* copy-on-write today — that is a real cost, but not one a wall-clock ratio
* can report reliably, as the numbers above show.
* Kotlin/Native (Apple and Linux) now backs LargeCache with StripedHashMap,
* O(1) per write by construction, so the guarantee does not need pinning
* there; on JVM/Android it comes from LargeCache being a concurrent map, and a
* regression in this class (someone reintroducing a copy-on-write map, or a
* per-publish full scan) shows up cleanly.
*/
class PoolEventOutboxScaleTest {
private val relay = NormalizedRelayUrl("wss://scale.relay.test")
@@ -21,8 +21,8 @@
package com.vitorpamplona.quartz.utils.cache
/**
* Linux/Native actual for [ConcurrentHashCache], over the same [StripedHashMap] as
* `LargeCache.linux.kt` — read its docs for why.
* Kotlin/Native actual (Apple and Linux) for [ConcurrentHashCache], over the same
* [StripedHashMap] as `LargeCache.native.kt` — read its docs for why.
*
* 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
@@ -21,7 +21,8 @@
package com.vitorpamplona.quartz.utils.cache
/**
* Linux/Native actual for [LargeCache] — the store behind Amethyst's `LocalCache`.
* Kotlin/Native actual (Apple and Linux) for [LargeCache] — the store behind Amethyst's
* `LocalCache`.
*
* All of the concurrency and the performance rationale lives in [StripedHashMap]; this
* class is only the [ICacheOperations] surface over it. The short version: `LocalCache`
@@ -40,8 +41,16 @@ package com.vitorpamplona.quartz.utils.cache
* - There is no snapshot and no defensive `entries.toList()`, so no
* `ConcurrentModificationException` window and no per-scan copy.
*
* Apple used to wrap charlietap's `CacheMap` instead. That is a left-right map: its
* `entries`/`keys`/`values` hand out live views of one of its two inner `HashMap`s and
* release the read guard before the caller iterates them, so every bulk operation here
* (`filter`, `map`, `keys()`, even a defensive `entries.toList()`) walked a `HashMap` a
* writer could be mutating. Kotlin/Native's `HashMap` detects that and throws
* `ConcurrentModificationException` — which is how `NostrClient`'s pool scans failed on
* the iOS simulator while a bunker round-trip subscribed and published from other threads.
*
* Iteration is weakly consistent and in bucket order. JVM/Android iterates in
* sorted-key order (`ConcurrentSkipListMap`) and Apple in hash order; nothing in the
* sorted-key order (`ConcurrentSkipListMap`); 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`.
@@ -43,7 +43,9 @@ import kotlin.concurrent.atomics.ExperimentalAtomicApi
* - **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).
* predicates reach back into the cache). Apple's former `CacheMap` (left-right over
* two `HashMap`s) skipped the copy and so let scans race writers into a
* `ConcurrentModificationException`.
*
* 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
@@ -27,11 +27,12 @@ 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.
* The native actual (Apple and Linux) 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 Linux copy-on-write version it replaced would fail every
* assertion here: its read-copy-write was not a CAS loop, so concurrent writers silently
* dropped each other's entries. The Apple `CacheMap` version it replaced failed
* [scansNeverThrowWhileKeysComeAndGo] with a `ConcurrentModificationException`.
*
* Uses `Worker` rather than coroutines on purpose — a coroutine dispatcher gives no
* guarantee of genuine parallelism, and parallelism is the whole point.
@@ -117,4 +118,37 @@ class LargeCacheConcurrencyTest {
assertEquals(1_000 + (workerCount / 2) * perWorker, cache.size())
}
@Test
fun scansNeverThrowWhileKeysComeAndGo() {
val cache = LargeCache<Int, Int>()
repeat(1_000) { cache.put(it, it) }
// The shape NostrClient's pool has: one thread adds and drops subscriptions while
// another walks the whole map to rebuild filters. Every bulk read must tolerate a
// structural change landing mid-walk, removals included.
inParallel { id ->
if (id % 2 == 0) {
repeat(perWorker) { i ->
val key = 1_000 + id * perWorker + i
cache.put(key, i)
cache.remove(key - 1)
}
} else {
repeat(200) {
cache.filter { _, v -> v >= 0 }
cache.map { k, _ -> k }
cache.mapNotNull { k, v -> if (v % 2 == 0) k else null }
cache.keys().forEach { check(it >= 0) }
cache.values().forEach { check(it >= 0) }
cache.forEach { k, _ -> check(k >= 0) }
}
}
}
// Each writer leaves only its last key behind; its first remove takes out 999
// (the first writer's) or a key that was never there (the others').
assertEquals(999 + workerCount / 2, cache.size())
assertEquals(998, cache.get(998))
}
}