mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 11:48:24 +00:00
perf(quartz): drop copy-on-write from linuxX64's LargeCache
LocalCache's linuxX64 store kept a LinkedHashMap inside an AtomicReference and replaced it wholesale on every write, so each put was O(n) in the size of the cache and filling it was O(n^2). It was not thread-safe either: the read-copy-write was not a CAS loop, so concurrent writers silently dropped each other's entries. Replace it with a mutable map guarded by PlatformLock plus a lazily rebuilt read snapshot. Point operations (get/put/remove/containsKey/size) are O(1) under the lock; bulk operations run against a point-in-time copy rebuilt at most once per write epoch, which also keeps caller-supplied lambdas out of the critical section — PlatformLock is not reentrant here and LocalCache predicates call back into the cache. Two behaviour fixes fall out of matching the JVM actual's putIfAbsent: createIfAbsent now reports true only when this call inserted (it previously returned get(key) != null, which also reported true when another thread had just created the entry), and getOrCreate publishes atomically. ConcurrentHashCache.linux gets the same treatment. Its only caller, CachingEventDecoder, writes once per event arriving from a relay, so the per-write map rebuild was the worst-placed copy of the three. None of this was caught because no CI job compiled or ran linuxX64. Add LargeCacheTest to commonTest as a cross-target contract for the ~40 methods each actual reimplements by hand, a linuxTest suite covering the concurrency this actual now has to get right on its own, and a CI leg that runs both on Linux Native. That leg is scoped to the cache and concurrency packages: the full linuxX64Test suite is 3,490 tests with 78 pre-existing failures, nearly all TODO() stubs in linux actuals that were never written (MLS crypto, the SQLite driver, NIP-44, Bolt12). Filling those in is its own project; the filter keeps the job meaningful and green, and widening it later is one line. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01HxQ1QuyzSkR38iFHbREjoS
This commit is contained in:
@@ -194,6 +194,76 @@ jobs:
|
||||
name: geode Test Reports
|
||||
path: geode/build/reports
|
||||
|
||||
# linuxX64 is the only target whose LargeCache / ConcurrentHashCache actuals are
|
||||
# hand-written concurrent maps rather than a delegation to a platform concurrent
|
||||
# collection, and until this job existed nothing ran them: the target was compiled
|
||||
# by no CI leg at all. That is how a copy-on-write LargeCache with O(n) writes and a
|
||||
# non-atomic read-copy-write (concurrent writers silently dropped entries) sat in
|
||||
# the tree unnoticed.
|
||||
#
|
||||
# Scoped to the cache/concurrency packages on purpose. The full linuxX64Test suite
|
||||
# is 3,490 tests with 78 pre-existing failures, essentially all of them `TODO()`
|
||||
# stubs in linux actuals that were never written — MLS crypto, the SQLite driver,
|
||||
# NIP-44, Bolt12 — plus two URL-handling divergences. Filling those in is its own
|
||||
# project; gating PRs on them today would just mean a permanently red job. The
|
||||
# filter keeps the leg meaningful and green, and widening it is a one-line change
|
||||
# once the native actuals land.
|
||||
#
|
||||
# This still compiles and links the whole module for linuxX64, so a commonMain or
|
||||
# commonTest source that reaches for a JVM-only API fails here too — on a target
|
||||
# with no Foundation to fall back on the way Apple has.
|
||||
test-quartz-linux-native:
|
||||
needs: lint
|
||||
runs-on: ubuntu-latest
|
||||
timeout-minutes: 45
|
||||
steps:
|
||||
- name: Checkout code
|
||||
uses: actions/checkout@v7
|
||||
|
||||
- name: Set up JDK 21
|
||||
uses: actions/setup-java@v6.0.0
|
||||
with:
|
||||
distribution: 'temurin'
|
||||
java-version: 21
|
||||
|
||||
- name: Set up Gradle
|
||||
uses: gradle/actions/setup-gradle@v6
|
||||
with:
|
||||
cache-read-only: ${{ github.ref != 'refs/heads/main' }}
|
||||
|
||||
# The Kotlin/Native toolchain (compiler distribution + LLVM + the sysroot) lands
|
||||
# in ~/.konan, which setup-gradle does not cache. Without this the job re-downloads
|
||||
# well over a gigabyte on every run. Keyed on the version catalog so a Kotlin bump
|
||||
# re-populates it.
|
||||
- name: Cache Kotlin/Native toolchain
|
||||
uses: actions/cache@v4
|
||||
with:
|
||||
path: ~/.konan
|
||||
key: konan-${{ runner.os }}-${{ hashFiles('gradle/libs.versions.toml') }}
|
||||
restore-keys: konan-${{ runner.os }}-
|
||||
|
||||
- name: Test Quartz caches on Linux Native
|
||||
run: |
|
||||
./gradlew :quartz:linuxX64Test \
|
||||
--tests "com.vitorpamplona.quartz.utils.cache.*" \
|
||||
--tests "com.vitorpamplona.quartz.utils.concurrent.*"
|
||||
|
||||
- name: Linux Native Test Report
|
||||
uses: mikepenz/action-junit-report@a9170d5795813c01ab4901ffb045b52bab4ab09d # v6.5.0
|
||||
if: always()
|
||||
with:
|
||||
report_paths: 'quartz/build/test-results/linuxX64Test/TEST-*.xml'
|
||||
annotate_only: true
|
||||
detailed_summary: true
|
||||
fail_on_failure: true
|
||||
|
||||
- name: Upload Linux Native Test Reports
|
||||
uses: actions/upload-artifact@v7
|
||||
if: failure()
|
||||
with:
|
||||
name: Quartz Linux Native Test Reports
|
||||
path: quartz/build/reports
|
||||
|
||||
test-quartz-ios:
|
||||
# Phase 1 of the iOS support plan
|
||||
# (amethyst/plans/2026-05-24-ios-support.md): keep :quartz green on iOS
|
||||
|
||||
+319
@@ -0,0 +1,319 @@
|
||||
/*
|
||||
* Copyright (c) 2025 Vitor Pamplona
|
||||
*
|
||||
* Permission is hereby granted, free of charge, to any person obtaining a copy of
|
||||
* this software and associated documentation files (the "Software"), to deal in
|
||||
* the Software without restriction, including without limitation the rights to use,
|
||||
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
|
||||
* Software, and to permit persons to whom the Software is furnished to do so,
|
||||
* subject to the following conditions:
|
||||
*
|
||||
* The above copyright notice and this permission notice shall be included in all
|
||||
* copies or substantial portions of the Software.
|
||||
*
|
||||
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
|
||||
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
|
||||
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
|
||||
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
|
||||
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
|
||||
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
||||
*/
|
||||
package com.vitorpamplona.quartz.utils.cache
|
||||
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFalse
|
||||
import kotlin.test.assertNull
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
/**
|
||||
* Cross-target contract for [LargeCache], the store behind `LocalCache`.
|
||||
*
|
||||
* There was no shared suite for this class: the JVM/Android actual is covered only
|
||||
* indirectly through `LocalCache`, and the linuxX64 and Apple actuals — hand-written
|
||||
* reimplementations of the same ~40 methods — were covered by nothing at all. This
|
||||
* runs the same assertions against whichever actual the target picked, so a divergence
|
||||
* shows up as a test failure instead of at runtime on one platform.
|
||||
*
|
||||
* Deliberately order-agnostic: iteration order is insertion order on linux, sorted-key
|
||||
* order on JVM/Android (`ConcurrentSkipListMap`) and hash order on Apple. Anything
|
||||
* order-sensitive is compared as a set or sorted first.
|
||||
*/
|
||||
class LargeCacheTest {
|
||||
private fun cacheOf(vararg pairs: Pair<String, Int>) =
|
||||
LargeCache<String, Int>().apply {
|
||||
pairs.forEach { put(it.first, it.second) }
|
||||
}
|
||||
|
||||
@Test
|
||||
fun emptyCache() {
|
||||
val cache = LargeCache<String, Int>()
|
||||
|
||||
assertEquals(0, cache.size())
|
||||
assertTrue(cache.isEmpty())
|
||||
assertNull(cache.get("a"))
|
||||
assertFalse(cache.containsKey("a"))
|
||||
assertTrue(cache.keys().isEmpty())
|
||||
assertTrue(cache.values().toList().isEmpty())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun putGetAndOverwrite() {
|
||||
val cache = cacheOf("a" to 1, "b" to 2)
|
||||
|
||||
assertEquals(1, cache.get("a"))
|
||||
assertEquals(2, cache.get("b"))
|
||||
assertEquals(2, cache.size())
|
||||
assertFalse(cache.isEmpty())
|
||||
assertTrue(cache.containsKey("a"))
|
||||
|
||||
cache.put("a", 10)
|
||||
|
||||
assertEquals(10, cache.get("a"))
|
||||
assertEquals(2, cache.size(), "overwriting a key must not grow the cache")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun removeReturnsOldValue() {
|
||||
val cache = cacheOf("a" to 1, "b" to 2)
|
||||
|
||||
assertEquals(1, cache.remove("a"))
|
||||
assertNull(cache.get("a"))
|
||||
assertFalse(cache.containsKey("a"))
|
||||
assertEquals(1, cache.size())
|
||||
|
||||
assertNull(cache.remove("a"), "removing an absent key returns null")
|
||||
assertEquals(1, cache.size())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun clearEmptiesEverything() {
|
||||
val cache = cacheOf("a" to 1, "b" to 2)
|
||||
|
||||
cache.clear()
|
||||
|
||||
assertEquals(0, cache.size())
|
||||
assertTrue(cache.isEmpty())
|
||||
assertNull(cache.get("a"))
|
||||
assertTrue(cache.keys().isEmpty())
|
||||
assertTrue(cache.values().toList().isEmpty())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun getOrCreateBuildsOnlyOnce() {
|
||||
val cache = LargeCache<String, Int>()
|
||||
var builds = 0
|
||||
|
||||
assertEquals(
|
||||
7,
|
||||
cache.getOrCreate("k") {
|
||||
builds++
|
||||
7
|
||||
},
|
||||
)
|
||||
assertEquals(
|
||||
7,
|
||||
cache.getOrCreate("k") {
|
||||
builds++
|
||||
99
|
||||
},
|
||||
)
|
||||
|
||||
assertEquals(1, builds, "the builder must not run for a key that is already present")
|
||||
assertEquals(7, cache.get("k"))
|
||||
assertEquals(1, cache.size())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun createIfAbsentReportsWhoInserted() {
|
||||
val cache = LargeCache<String, Int>()
|
||||
|
||||
assertTrue(cache.createIfAbsent("k") { 1 }, "the first call inserts")
|
||||
assertFalse(cache.createIfAbsent("k") { 2 }, "the second call must not report an insert")
|
||||
|
||||
assertEquals(1, cache.get("k"), "a losing createIfAbsent must not overwrite")
|
||||
assertEquals(1, cache.size())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun keysAndValuesSeeLaterWrites() {
|
||||
val cache = cacheOf("a" to 1)
|
||||
|
||||
// Reads the collections first, so an implementation that caches a snapshot has
|
||||
// one to go stale.
|
||||
assertEquals(setOf("a"), cache.keys().toSet())
|
||||
assertEquals(listOf(1), cache.values().toList())
|
||||
|
||||
cache.put("b", 2)
|
||||
|
||||
assertEquals(setOf("a", "b"), cache.keys().toSet())
|
||||
assertEquals(listOf(1, 2), cache.values().sorted())
|
||||
|
||||
cache.remove("a")
|
||||
|
||||
assertEquals(setOf("b"), cache.keys().toSet())
|
||||
assertEquals(listOf(2), cache.values().toList())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun bulkReadsSeeLaterWrites() {
|
||||
val cache = cacheOf("a" to 1)
|
||||
|
||||
// Same idea for the collector path: warm every kind of bulk read, then mutate
|
||||
// and re-read. Guards the lazily rebuilt snapshot in the linux actual.
|
||||
assertEquals(1, cache.count { _, _ -> true })
|
||||
assertEquals(1, cache.sumOf { _, v -> v })
|
||||
|
||||
cache.put("b", 2)
|
||||
assertEquals(2, cache.count { _, _ -> true })
|
||||
assertEquals(3, cache.sumOf { _, v -> v })
|
||||
|
||||
cache.put("a", 10)
|
||||
assertEquals(12, cache.sumOf { _, v -> v }, "an overwrite must invalidate a cached read view")
|
||||
|
||||
cache.remove("b")
|
||||
assertEquals(10, cache.sumOf { _, v -> v })
|
||||
|
||||
cache.getOrCreate("c") { 5 }
|
||||
assertEquals(15, cache.sumOf { _, v -> v })
|
||||
|
||||
cache.createIfAbsent("d") { 100 }
|
||||
assertEquals(115, cache.sumOf { _, v -> v })
|
||||
|
||||
cache.clear()
|
||||
assertEquals(0, cache.count { _, _ -> true })
|
||||
assertEquals(0, cache.sumOf { _, v -> v })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun forEachVisitsEveryEntry() {
|
||||
val cache = cacheOf("a" to 1, "b" to 2, "c" to 3)
|
||||
|
||||
val seen = mutableMapOf<String, Int>()
|
||||
cache.forEach { k, v -> seen[k] = v }
|
||||
|
||||
assertEquals(mapOf("a" to 1, "b" to 2, "c" to 3), seen)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun filterAndCount() {
|
||||
val cache = cacheOf("a" to 1, "b" to 2, "c" to 3, "d" to 4)
|
||||
|
||||
assertEquals(listOf(2, 4), cache.filter { _, v -> v % 2 == 0 }.sorted())
|
||||
assertEquals(setOf(2, 4), cache.filterIntoSet { _, v -> v % 2 == 0 })
|
||||
assertEquals(2, cache.count { _, v -> v % 2 == 0 })
|
||||
assertEquals(1, cache.count { k, _ -> k == "a" })
|
||||
assertEquals(0, cache.count { _, _ -> false })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun mapVariants() {
|
||||
val cache = cacheOf("a" to 1, "b" to 2, "c" to 3)
|
||||
|
||||
assertEquals(listOf(2, 4, 6), cache.map { _, v -> v * 2 }.sorted())
|
||||
assertEquals(listOf(2, 3), cache.mapNotNull { _, v -> if (v > 1) v else null }.sorted())
|
||||
assertEquals(setOf(2, 3), cache.mapNotNullIntoSet { _, v -> if (v > 1) v else null })
|
||||
assertEquals(
|
||||
listOf(1, 1, 2, 2, 3, 3),
|
||||
cache.mapFlatten { _, v -> listOf(v, v) }.sorted(),
|
||||
)
|
||||
assertEquals(
|
||||
setOf(1, 2, 3),
|
||||
cache.mapFlattenIntoSet { _, v -> listOf(v, v) },
|
||||
)
|
||||
assertEquals(
|
||||
listOf(2, 3),
|
||||
cache.mapFlatten { _, v -> if (v > 1) listOf(v) else null }.sorted(),
|
||||
"a null collection contributes nothing",
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun aggregates() {
|
||||
val cache = cacheOf("a" to 1, "b" to 2, "c" to 3, "d" to 4)
|
||||
|
||||
assertEquals(10, cache.sumOf { _, v -> v })
|
||||
assertEquals(10L, cache.sumOfLong { _, v -> v.toLong() })
|
||||
assertEquals(4, cache.maxOrNullOf({ _, _ -> true }, naturalOrder<Int>()))
|
||||
assertEquals(3, cache.maxOrNullOf({ _, v -> v % 2 == 1 }, naturalOrder<Int>()))
|
||||
assertNull(cache.maxOrNullOf({ _, _ -> false }, naturalOrder<Int>()))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun groupings() {
|
||||
val cache = cacheOf("a" to 1, "b" to 2, "c" to 3, "d" to 4)
|
||||
|
||||
assertEquals(
|
||||
mapOf(0 to listOf(2, 4), 1 to listOf(1, 3)),
|
||||
cache.groupBy<Int> { _, v -> v % 2 }.mapValues { it.value.sorted() },
|
||||
)
|
||||
assertEquals(mapOf(0 to 2, 1 to 2), cache.countByGroup<Int> { _, v -> v % 2 })
|
||||
assertEquals(
|
||||
mapOf(0 to 6L, 1 to 4L),
|
||||
cache.sumByGroup({ _, v -> v % 2 }, { _, v -> v.toLong() }),
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun associates() {
|
||||
val cache = cacheOf("a" to 1, "b" to 2)
|
||||
|
||||
assertEquals(mapOf(1 to "a", 2 to "b"), cache.associate { k, v -> v to k })
|
||||
assertEquals(mapOf("a" to 2, "b" to 4), cache.associateWith { _, v -> v * 2 })
|
||||
assertEquals(mapOf("a" to null, "b" to 4), cache.associateWith { _, v -> if (v > 1) v * 2 else null })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun joinToStringRendersEveryEntry() {
|
||||
val cache = cacheOf("a" to 1, "b" to 2, "c" to 3)
|
||||
|
||||
val rendered =
|
||||
cache.joinToString(
|
||||
separator = ",",
|
||||
prefix = "[",
|
||||
postfix = "]",
|
||||
limit = -1,
|
||||
truncated = "...",
|
||||
) { k, v -> "$k=$v" }
|
||||
|
||||
assertTrue(rendered.startsWith("[") && rendered.endsWith("]"), rendered)
|
||||
assertEquals(
|
||||
listOf("a=1", "b=2", "c=3"),
|
||||
rendered
|
||||
.removePrefix("[")
|
||||
.removeSuffix("]")
|
||||
.split(",")
|
||||
.sorted(),
|
||||
)
|
||||
|
||||
assertEquals("[]", LargeCache<String, Int>().joinToString(",", "[", "]", -1, "...") { k, v -> "$k=$v" })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun insertingManyEntriesStaysLinear() {
|
||||
// Regression guard for the copy-on-write linux actual this replaced, where every
|
||||
// put rebuilt the whole map: this loop cost ~1.25 billion entry copies there and
|
||||
// is instant against any O(1)-write implementation. No wall-clock assertion —
|
||||
// the run time itself is the signal.
|
||||
val n = 50_000
|
||||
val cache = LargeCache<Int, Int>()
|
||||
|
||||
for (i in 0 until n) {
|
||||
cache.put(i, i)
|
||||
// A full scan every thousandth insert, so an implementation that rebuilds a
|
||||
// read view on write has to rebuild it ~50 times rather than amortize it away.
|
||||
if (i % 1_000 == 0) assertEquals(i + 1, cache.count { _, _ -> true })
|
||||
}
|
||||
|
||||
assertEquals(n, cache.size())
|
||||
assertEquals(0, cache.get(0))
|
||||
assertEquals(n - 1, cache.get(n - 1))
|
||||
assertEquals(n.toLong() * (n - 1) / 2, cache.sumOfLong { _, v -> v.toLong() })
|
||||
|
||||
for (i in 0 until n step 2) cache.remove(i)
|
||||
|
||||
assertEquals(n / 2, cache.size())
|
||||
assertNull(cache.get(0))
|
||||
assertEquals(1, cache.get(1))
|
||||
}
|
||||
}
|
||||
Vendored
+20
-13
@@ -20,30 +20,37 @@
|
||||
*/
|
||||
package com.vitorpamplona.quartz.utils.cache
|
||||
|
||||
import kotlin.concurrent.AtomicReference
|
||||
import com.vitorpamplona.quartz.utils.concurrent.PlatformLock
|
||||
import com.vitorpamplona.quartz.utils.concurrent.withLock
|
||||
|
||||
// Copy-on-write, mirroring LargeCache.linux: correct and simple; the linux
|
||||
// target is CI-only so write cost is acceptable.
|
||||
/**
|
||||
* Linux/Native actual for [ConcurrentHashCache].
|
||||
*
|
||||
* Was copy-on-write, mirroring the old `LargeCache.linux`: every [put] rebuilt the
|
||||
* whole map under a CAS retry loop, so writes were O(n) and a decode burst against a
|
||||
* warm cache was O(n^2). That is a bad shape for this class in particular — its only
|
||||
* caller, `CachingEventDecoder`, writes once per event arriving from a relay.
|
||||
*
|
||||
* Now a plain [HashMap] guarded by a [PlatformLock]: O(1) writes, and no CAS retry
|
||||
* to livelock under a write burst. Same lock choice and the same residual (a spin
|
||||
* lock on this target) as `LargeCache.linux.kt` — see its docs.
|
||||
*/
|
||||
actual class ConcurrentHashCache<K : Any, V : Any> {
|
||||
private val mapRef = AtomicReference(HashMap<K, V>())
|
||||
private val lock = PlatformLock()
|
||||
private val map = HashMap<K, V>()
|
||||
|
||||
actual fun get(key: K): V? = mapRef.value[key]
|
||||
actual fun get(key: K): V? = lock.withLock { map[key] }
|
||||
|
||||
actual fun put(
|
||||
key: K,
|
||||
value: V,
|
||||
) {
|
||||
while (true) {
|
||||
val current = mapRef.value
|
||||
val copy = HashMap(current)
|
||||
copy[key] = value
|
||||
if (mapRef.compareAndSet(current, copy)) return
|
||||
}
|
||||
lock.withLock { map[key] = value }
|
||||
}
|
||||
|
||||
actual fun size(): Int = mapRef.value.size
|
||||
actual fun size(): Int = lock.withLock { map.size }
|
||||
|
||||
actual fun clear() {
|
||||
mapRef.value = HashMap()
|
||||
lock.withLock { map.clear() }
|
||||
}
|
||||
}
|
||||
|
||||
+134
-35
@@ -20,41 +20,112 @@
|
||||
*/
|
||||
package com.vitorpamplona.quartz.utils.cache
|
||||
|
||||
import kotlin.concurrent.AtomicReference
|
||||
import com.vitorpamplona.quartz.utils.concurrent.PlatformLock
|
||||
import com.vitorpamplona.quartz.utils.concurrent.withLock
|
||||
|
||||
/**
|
||||
* Linux/Native actual for [LargeCache] — the store behind Amethyst's `LocalCache`.
|
||||
*
|
||||
* ## Why this is not copy-on-write any more
|
||||
*
|
||||
* The first cut of this file kept a `LinkedHashMap` inside an `AtomicReference` and
|
||||
* replaced it wholesale on every write. That made each [put] **O(n) in the size of
|
||||
* the cache**: inserting n entries cost O(n^2) copies, so a cache holding 100k notes
|
||||
* paid a 100k-entry map copy per arriving event. It was also *not* actually
|
||||
* thread-safe — the read-copy-write was not a CAS loop, so two concurrent writers
|
||||
* silently dropped one of the two writes.
|
||||
*
|
||||
* That went unnoticed because no CI job runs the linuxX64 target, so nothing ever
|
||||
* pushed volume through this class.
|
||||
*
|
||||
* ## What it does instead
|
||||
*
|
||||
* A single mutable [LinkedHashMap] guarded by a [PlatformLock], plus a lazily built
|
||||
* read snapshot:
|
||||
*
|
||||
* - **Point operations** ([get], [put], [remove], [containsKey], [size], …) take the
|
||||
* lock, touch the live map, and return. O(1), no copying.
|
||||
* - **Bulk operations** (`filter`/`map`/`forEach`/…) run against [cachedSnapshot], a
|
||||
* point-in-time copy rebuilt on the first bulk call after a write and reused until
|
||||
* the next write. So a copy costs O(n) at most once per write epoch, on operations
|
||||
* that are already O(n), and a run of reads with no interleaved write copies
|
||||
* nothing at all.
|
||||
*
|
||||
* The snapshot is what lets the bulk operations invoke caller-supplied lambdas
|
||||
* **outside** the critical section. That matters: `PlatformLock` is not reentrant on
|
||||
* this target, and `LocalCache` predicates routinely call back into the same cache —
|
||||
* running them under the lock would self-deadlock. It also removes the
|
||||
* `ConcurrentModificationException` window the previous `entries.toList()` dance was
|
||||
* working around.
|
||||
*
|
||||
* ## Known residual
|
||||
*
|
||||
* Building the snapshot holds the lock for O(n), and the linux `PlatformLock` is a
|
||||
* spin lock (Kotlin/Native ships no parking lock and there is no Foundation here —
|
||||
* see `PlatformLock.linux.kt`). A writer racing a snapshot build therefore busy-waits
|
||||
* for the duration of the copy. That is still strictly better than what it replaces —
|
||||
* copy-on-write did the same O(n) copy on *every write* and lost concurrent ones —
|
||||
* but it is the reason `PlatformLock.linux.kt` flags a pthread mutex as the next step
|
||||
* if this target ever hosts a genuinely contended workload.
|
||||
*
|
||||
* Ordering note: iteration follows insertion order here, sorted-key order on
|
||||
* JVM/Android (`ConcurrentSkipListMap`) and hash order on Apple. Nothing in the
|
||||
* codebase depends on a specific order, and the `from`/`to` range overloads below
|
||||
* degrade to a full scan on this target exactly as they did before — they have no
|
||||
* callers outside the JVM-only `LargeSoftCache`.
|
||||
*/
|
||||
actual class LargeCache<K, V> : ICacheOperations<K, V> {
|
||||
private val mapRef = AtomicReference(LinkedHashMap<K, V>())
|
||||
private val lock = PlatformLock()
|
||||
|
||||
private inline fun <R> withMap(block: (LinkedHashMap<K, V>) -> R): R = block(mapRef.value)
|
||||
/** The live store. Every access must hold [lock]. */
|
||||
private val map = LinkedHashMap<K, V>()
|
||||
|
||||
private inline fun mutate(block: (LinkedHashMap<K, V>) -> Unit) {
|
||||
val copy = LinkedHashMap(mapRef.value)
|
||||
block(copy)
|
||||
mapRef.value = copy
|
||||
}
|
||||
/**
|
||||
* A copy of [map] handed to bulk operations so they can run caller lambdas
|
||||
* without holding [lock]. Null means "stale, rebuild on next use". Guarded by
|
||||
* [lock]; never mutated once published, so readers may keep it as long as they
|
||||
* like.
|
||||
*/
|
||||
private var cachedSnapshot: Map<K, V>? = null
|
||||
|
||||
actual fun keys(): Set<K> = withMap { LinkedHashSet(it.keys) }
|
||||
|
||||
actual fun values(): Iterable<V> = withMap { ArrayList(it.values) }
|
||||
|
||||
actual fun get(key: K): V? = withMap { it[key] }
|
||||
|
||||
actual fun remove(key: K): V? {
|
||||
val current = mapRef.value
|
||||
val value = current[key]
|
||||
if (value != null) {
|
||||
mutate { it.remove(key) }
|
||||
private fun snapshot(): Map<K, V> =
|
||||
lock.withLock {
|
||||
cachedSnapshot ?: LinkedHashMap(map).also { cachedSnapshot = it }
|
||||
}
|
||||
return value
|
||||
}
|
||||
|
||||
actual fun isEmpty(): Boolean = withMap { it.isEmpty() }
|
||||
/** Runs [block] over a stable snapshot, outside the lock. */
|
||||
private inline fun <R> withMap(block: (Map<K, V>) -> R): R = block(snapshot())
|
||||
|
||||
/** Runs [block] over the live map under [lock] without invalidating the snapshot. */
|
||||
private inline fun <R> read(block: (MutableMap<K, V>) -> R): R = lock.withLock { block(map) }
|
||||
|
||||
/** Runs [block] over the live map under [lock] and drops the read snapshot. */
|
||||
private inline fun <R> mutate(block: (MutableMap<K, V>) -> R): R =
|
||||
lock.withLock {
|
||||
cachedSnapshot = null
|
||||
block(map)
|
||||
}
|
||||
|
||||
actual fun keys(): Set<K> = snapshot().keys
|
||||
|
||||
actual fun values(): Iterable<V> = snapshot().values
|
||||
|
||||
actual fun get(key: K): V? = read { it[key] }
|
||||
|
||||
actual fun remove(key: K): V? =
|
||||
lock.withLock {
|
||||
val removed = map.remove(key)
|
||||
if (removed != null) cachedSnapshot = null
|
||||
removed
|
||||
}
|
||||
|
||||
actual fun isEmpty(): Boolean = read { it.isEmpty() }
|
||||
|
||||
actual fun clear() {
|
||||
mapRef.value = LinkedHashMap()
|
||||
mutate { it.clear() }
|
||||
}
|
||||
|
||||
actual fun containsKey(key: K): Boolean = withMap { it.containsKey(key) }
|
||||
actual fun containsKey(key: K): Boolean = read { it.containsKey(key) }
|
||||
|
||||
actual fun put(
|
||||
key: K,
|
||||
@@ -63,33 +134,61 @@ actual class LargeCache<K, V> : ICacheOperations<K, V> {
|
||||
mutate { it[key] = value }
|
||||
}
|
||||
|
||||
/**
|
||||
* Mirrors the JVM actual's `putIfAbsent`: [builder] runs outside the lock (it is
|
||||
* caller code and must not be able to re-enter a non-reentrant lock), and the
|
||||
* insert is only published if no one won the race in the meantime.
|
||||
*/
|
||||
actual fun getOrCreate(
|
||||
key: K,
|
||||
builder: (key: K) -> V,
|
||||
): V {
|
||||
val existing = get(key)
|
||||
if (existing != null) return existing
|
||||
read { it[key] }?.let { return it }
|
||||
|
||||
val newObject = builder(key)
|
||||
mutate { it[key] = newObject }
|
||||
return get(key) ?: newObject
|
||||
|
||||
return lock.withLock {
|
||||
val existing = map[key]
|
||||
if (existing != null) {
|
||||
existing
|
||||
} else {
|
||||
map[key] = newObject
|
||||
cachedSnapshot = null
|
||||
newObject
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* True only when *this* call inserted the value — matching the JVM actual's
|
||||
* `putIfAbsent(key, newObject) == null`. The previous implementation returned
|
||||
* `get(key) != null`, which also reported true when another thread had just
|
||||
* created the entry, double-firing whatever the caller does with a fresh key.
|
||||
*/
|
||||
actual fun createIfAbsent(
|
||||
key: K,
|
||||
builder: (key: K) -> V,
|
||||
): Boolean {
|
||||
val existing = get(key)
|
||||
if (existing != null) return false
|
||||
if (read { it.containsKey(key) }) return false
|
||||
|
||||
val newObject = builder(key)
|
||||
mutate { it[key] = newObject }
|
||||
return get(key) != null
|
||||
|
||||
return lock.withLock {
|
||||
if (map.containsKey(key)) {
|
||||
false
|
||||
} else {
|
||||
map[key] = newObject
|
||||
cachedSnapshot = null
|
||||
true
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
actual override fun size(): Int = withMap { it.size }
|
||||
actual override fun size(): Int = read { it.size }
|
||||
|
||||
actual override fun forEach(consumer: ICacheBiConsumer<K, V>) {
|
||||
// Snapshot entries to avoid ConcurrentModificationException
|
||||
withMap { map -> map.entries.toList() }.forEach { consumer.accept(it.key, it.value) }
|
||||
// The snapshot is already immutable, so no defensive entries.toList() is needed.
|
||||
withMap { map -> map.forEach { consumer.accept(it.key, it.value) } }
|
||||
}
|
||||
|
||||
actual override fun filter(consumer: CacheCollectors.BiFilter<K, V>): List<V> = withMap { map -> map.filter { consumer.filter(it.key, it.value) }.values.toList() }
|
||||
|
||||
Vendored
+120
@@ -0,0 +1,120 @@
|
||||
/*
|
||||
* Copyright (c) 2025 Vitor Pamplona
|
||||
*
|
||||
* Permission is hereby granted, free of charge, to any person obtaining a copy of
|
||||
* this software and associated documentation files (the "Software"), to deal in
|
||||
* the Software without restriction, including without limitation the rights to use,
|
||||
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
|
||||
* Software, and to permit persons to whom the Software is furnished to do so,
|
||||
* subject to the following conditions:
|
||||
*
|
||||
* The above copyright notice and this permission notice shall be included in all
|
||||
* copies or substantial portions of the Software.
|
||||
*
|
||||
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
|
||||
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
|
||||
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
|
||||
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
|
||||
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
|
||||
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
||||
*/
|
||||
package com.vitorpamplona.quartz.utils.cache
|
||||
|
||||
import kotlin.native.concurrent.ObsoleteWorkersApi
|
||||
import kotlin.native.concurrent.TransferMode
|
||||
import kotlin.native.concurrent.Worker
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
|
||||
/**
|
||||
* The linux actual is the only [LargeCache] whose thread safety is hand-rolled rather
|
||||
* than delegated to a concurrent map, so it gets its own multi-threaded test. The
|
||||
* copy-on-write version this replaced would fail every assertion here: its
|
||||
* read-copy-write was not a CAS loop, so concurrent writers silently dropped each
|
||||
* other's entries.
|
||||
*
|
||||
* Uses `Worker` rather than coroutines on purpose — a coroutine dispatcher gives no
|
||||
* guarantee of genuine parallelism, and parallelism is the whole point.
|
||||
*/
|
||||
@OptIn(ObsoleteWorkersApi::class)
|
||||
class LargeCacheConcurrencyTest {
|
||||
private val workerCount = 4
|
||||
private val perWorker = 5_000
|
||||
|
||||
private fun <R> inParallel(job: (workerId: Int) -> R): List<R> {
|
||||
val workers = List(workerCount) { Worker.start() }
|
||||
val futures =
|
||||
workers.mapIndexed { id, worker ->
|
||||
worker.execute(TransferMode.SAFE, { Pair(job, id) }) { (block, workerId) -> block(workerId) }
|
||||
}
|
||||
val results = futures.map { it.result }
|
||||
workers.forEach { it.requestTermination().result }
|
||||
return results
|
||||
}
|
||||
|
||||
@Test
|
||||
fun concurrentPutsKeepEveryEntry() {
|
||||
val cache = LargeCache<Int, Int>()
|
||||
|
||||
inParallel { id ->
|
||||
repeat(perWorker) { i -> cache.put(id * perWorker + i, i) }
|
||||
}
|
||||
|
||||
assertEquals(workerCount * perWorker, cache.size())
|
||||
for (id in 0 until workerCount) {
|
||||
assertEquals(0, cache.get(id * perWorker))
|
||||
assertEquals(perWorker - 1, cache.get(id * perWorker + perWorker - 1))
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun concurrentGetOrCreateBuildsOneValuePerKey() {
|
||||
val cache = LargeCache<Int, String>()
|
||||
|
||||
// Every worker races for the same 500 keys. Whoever wins, all of them must end
|
||||
// up holding the identical instance, and the cache must hold exactly 500.
|
||||
val seen =
|
||||
inParallel { _ ->
|
||||
(0 until 500).map { key -> cache.getOrCreate(key) { "v$it" } }
|
||||
}
|
||||
|
||||
assertEquals(500, cache.size())
|
||||
seen.forEach { perWorkerValues ->
|
||||
perWorkerValues.forEachIndexed { key, value ->
|
||||
assertEquals(cache.get(key), value, "getOrCreate handed out a value it did not publish")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun concurrentCreateIfAbsentReportsExactlyOneInsertPerKey() {
|
||||
val cache = LargeCache<Int, Int>()
|
||||
|
||||
val insertsPerWorker = inParallel { _ -> (0 until 500).count { key -> cache.createIfAbsent(key) { key } } }
|
||||
|
||||
assertEquals(500, cache.size())
|
||||
assertEquals(500, insertsPerWorker.sum(), "exactly one caller per key may report an insert")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun bulkReadsStayConsistentWhileWritesLand() {
|
||||
val cache = LargeCache<Int, Int>()
|
||||
repeat(1_000) { cache.put(it, it) }
|
||||
|
||||
// Writers churn the map while readers walk snapshots of it. A reader must never
|
||||
// see a torn map, and must never crash on a concurrent modification.
|
||||
inParallel { id ->
|
||||
if (id % 2 == 0) {
|
||||
repeat(perWorker) { i -> cache.put(1_000 + id * perWorker + i, i) }
|
||||
} else {
|
||||
repeat(200) {
|
||||
val sum = cache.sumOfLong { _, v -> v.toLong() }
|
||||
check(sum >= 0) { "unexpected negative sum $sum" }
|
||||
cache.count { _, _ -> true }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
assertEquals(1_000 + (workerCount / 2) * perWorker, cache.size())
|
||||
}
|
||||
}
|
||||
Vendored
+73
@@ -0,0 +1,73 @@
|
||||
/*
|
||||
* Copyright (c) 2025 Vitor Pamplona
|
||||
*
|
||||
* Permission is hereby granted, free of charge, to any person obtaining a copy of
|
||||
* this software and associated documentation files (the "Software"), to deal in
|
||||
* the Software without restriction, including without limitation the rights to use,
|
||||
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
|
||||
* Software, and to permit persons to whom the Software is furnished to do so,
|
||||
* subject to the following conditions:
|
||||
*
|
||||
* The above copyright notice and this permission notice shall be included in all
|
||||
* copies or substantial portions of the Software.
|
||||
*
|
||||
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
|
||||
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
|
||||
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
|
||||
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
|
||||
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
|
||||
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
||||
*/
|
||||
package com.vitorpamplona.quartz.utils.cache
|
||||
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertContentEquals
|
||||
import kotlin.test.assertEquals
|
||||
|
||||
/**
|
||||
* The `from`/`to` overloads narrow to a key range only where the backing map is sorted
|
||||
* (`ConcurrentSkipListMap` on JVM/Android). This target's map is not sorted, and `K`
|
||||
* carries no `Comparable` bound to sort it by, so every range overload deliberately
|
||||
* degrades to a full scan and leaves the caller's predicate to do the filtering.
|
||||
*
|
||||
* That is safe today because the range overloads have no callers outside the JVM-only
|
||||
* `LargeSoftCache`. This test pins the behaviour so the fallback stays *consistent*
|
||||
* with the unbounded form — the property anything sharing code across targets would
|
||||
* rely on — rather than silently returning a different set on this platform.
|
||||
*
|
||||
* Linux-only on purpose: it is a statement about this actual. The Apple actual walks
|
||||
* its map by index instead, which is a different (and order-dependent) contract.
|
||||
*/
|
||||
class LargeCacheRangeFallbackTest {
|
||||
private val cache =
|
||||
LargeCache<String, Int>().apply {
|
||||
put("a", 1)
|
||||
put("b", 2)
|
||||
put("c", 3)
|
||||
}
|
||||
|
||||
private val all = CacheCollectors.BiFilter<String, Int> { _, _ -> true }
|
||||
|
||||
@Test
|
||||
fun rangeOverloadsMatchTheirUnboundedForm() {
|
||||
assertContentEquals(cache.filter(all), cache.filter("a", "c", all))
|
||||
assertEquals(cache.filterIntoSet(all), cache.filterIntoSet("a", "c", all))
|
||||
assertEquals(cache.count(all), cache.count("a", "c", all))
|
||||
assertEquals(cache.sumOf { _, v -> v }, cache.sumOf("a", "c") { _, v -> v })
|
||||
assertEquals(cache.sumOfLong { _, v -> v.toLong() }, cache.sumOfLong("a", "c") { _, v -> v.toLong() })
|
||||
assertContentEquals(cache.map<Int> { _, v -> v }, cache.map<Int>("a", "c") { _, v -> v })
|
||||
assertContentEquals(cache.mapNotNull<Int> { _, v -> v }, cache.mapNotNull<Int>("a", "c") { _, v -> v })
|
||||
assertEquals(cache.associate { k, v -> k to v }, cache.associate("a", "c") { k, v -> k to v })
|
||||
assertEquals(cache.associateWith { _, v -> v }, cache.associateWith("a", "c") { _, v -> v })
|
||||
assertEquals(cache.countByGroup<Int> { _, v -> v % 2 }, cache.countByGroup<Int>("a", "c") { _, v -> v % 2 })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun aNarrowerRangeStillScansEverything() {
|
||||
// Documents the degradation explicitly: on a sorted target this would return
|
||||
// only "a", here it returns all three. Anyone who later gives this actual real
|
||||
// range support should update this test rather than discover the difference in
|
||||
// production.
|
||||
assertEquals(3, cache.count("a", "a", all))
|
||||
}
|
||||
}
|
||||
+6
-4
@@ -23,10 +23,12 @@ package com.vitorpamplona.quartz.utils.concurrent
|
||||
import kotlin.concurrent.atomics.AtomicReference
|
||||
import kotlin.concurrent.atomics.ExperimentalAtomicApi
|
||||
|
||||
// Copy-on-write, mirroring ConcurrentHashCache.linux: correct and simple. The
|
||||
// native targets never run the crawl this backs (it is JVM/Android-only work);
|
||||
// they only compile it, so the O(n)-per-write cost is irrelevant. A CAS retry
|
||||
// loop gives getOrPut/merge the same atomicity the JVM actual gets for free.
|
||||
// Copy-on-write: correct and simple. Unlike LargeCache/ConcurrentHashCache — which
|
||||
// back LocalCache and the event decoder and were moved off copy-on-write for exactly
|
||||
// this reason — the native targets never run the crawl this backs (it is
|
||||
// JVM/Android-only work); they only compile it, so the O(n)-per-write cost is
|
||||
// irrelevant. A CAS retry loop gives getOrPut/merge the same atomicity the JVM actual
|
||||
// gets for free.
|
||||
@OptIn(ExperimentalAtomicApi::class)
|
||||
actual class ConcurrentMap<K : Any, V : Any> {
|
||||
private val ref = AtomicReference(HashMap<K, V>())
|
||||
|
||||
Reference in New Issue
Block a user