From 5e662fef0b908f6e33f9d837328588766de4f106 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 1 Sep 2026 18:36:33 +0000 Subject: [PATCH] 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))) + } +}