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))) + } +}