mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 03:38:23 +00:00
perf: add lock-free-read ConcurrentLruCache, use on two hot read paths
Part C of the dispatchers/thread-caps audit. Both LnurlEndpointCache and DesktopCachedRichTextParser were bounded caches backed by a LinkedHashMap behind a single monitor (@Synchronized / Collections.synchronizedMap with accessOrder). An access-order map structurally mutates on get, so every read took the lock — serializing all readers on paths that are hot (kind-9735 zap-receipt validation; feed rich-text rendering). Add ConcurrentLruCache<K, V> in quartz utils: storage is a ConcurrentHashMap so get is lock-free; writes + eviction run under a small write lock that is off the read path. Eviction is least-recently-put order (get does not refresh recency) — exactly what LnurlEndpointCache already did, and fine for the deterministic rich-text parse cache. Point both caches at the shared helper. Covered by a new ConcurrentLruCacheTest (round-trip, eviction order, re-put recency refresh, get-does-not-refresh, clear, and a concurrent size-bound smoke test); the existing LnurlEndpointCacheTest still passes unchanged. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01ANuUziXKRafSTBxbh4SMoq
This commit is contained in:
+6
-8
@@ -23,25 +23,23 @@ package com.vitorpamplona.amethyst.desktop.service
|
||||
import com.vitorpamplona.amethyst.commons.model.ImmutableListOfLists
|
||||
import com.vitorpamplona.amethyst.commons.richtext.RichTextParser
|
||||
import com.vitorpamplona.amethyst.commons.richtext.RichTextViewerState
|
||||
import com.vitorpamplona.quartz.utils.cache.ConcurrentLruCache
|
||||
|
||||
object DesktopCachedRichTextParser {
|
||||
private const val MAX_CACHE_SIZE = 50
|
||||
|
||||
private val cache =
|
||||
java.util.Collections.synchronizedMap(
|
||||
object : LinkedHashMap<String, RichTextViewerState>(64, 0.75f, true) {
|
||||
override fun removeEldestEntry(eldest: Map.Entry<String, RichTextViewerState>) = size > MAX_CACHE_SIZE
|
||||
},
|
||||
)
|
||||
// Lock-free get on the feed rich-text render path; the previous access-order
|
||||
// synchronizedMap took a monitor even on reads.
|
||||
private val cache = ConcurrentLruCache<String, RichTextViewerState>(MAX_CACHE_SIZE)
|
||||
|
||||
fun parseText(
|
||||
content: String,
|
||||
tags: ImmutableListOfLists<String>,
|
||||
callbackUri: String? = null,
|
||||
): RichTextViewerState {
|
||||
cache[content]?.let { return it }
|
||||
cache.get(content)?.let { return it }
|
||||
val state = RichTextParser().parseText(content, tags, callbackUri)
|
||||
cache[content] = state
|
||||
cache.put(content, state)
|
||||
return state
|
||||
}
|
||||
|
||||
|
||||
+9
-21
@@ -20,6 +20,8 @@
|
||||
*/
|
||||
package com.vitorpamplona.quartz.nip57Zaps.validate
|
||||
|
||||
import com.vitorpamplona.quartz.utils.cache.ConcurrentLruCache
|
||||
|
||||
/**
|
||||
* Process-wide cache of LNURL-pay endpoint metadata, keyed by the canonical
|
||||
* `/.well-known/lnurlp/<user>` URL the recipient resolves to.
|
||||
@@ -37,37 +39,23 @@ package com.vitorpamplona.quartz.nip57Zaps.validate
|
||||
object LnurlEndpointCache {
|
||||
private const val MAX_ENTRIES = 1000
|
||||
|
||||
// Insertion-ordered map so we can evict the oldest entry once we hit the cap.
|
||||
// Synchronized externally — every mutating call holds the monitor.
|
||||
private val cache: LinkedHashMap<String, LnurlEndpointInfo> = LinkedHashMap()
|
||||
// Bounded cache with a lock-free get — hot on the zap-validation read path.
|
||||
// Eviction is least-recently-put (a get does not refresh recency), matching
|
||||
// the previous LinkedHashMap-based behaviour where only put reordered.
|
||||
private val cache = ConcurrentLruCache<String, LnurlEndpointInfo>(MAX_ENTRIES)
|
||||
|
||||
@Synchronized
|
||||
fun get(url: String): LnurlEndpointInfo? = cache[LnurlForm.normalizeUrl(url)]
|
||||
fun get(url: String): LnurlEndpointInfo? = cache.get(LnurlForm.normalizeUrl(url))
|
||||
|
||||
@Synchronized
|
||||
fun put(
|
||||
url: String,
|
||||
info: LnurlEndpointInfo,
|
||||
) {
|
||||
val key = LnurlForm.normalizeUrl(url)
|
||||
// Re-insert so the entry becomes "youngest" in iteration order.
|
||||
cache.remove(key)
|
||||
cache[key] = info
|
||||
if (cache.size > MAX_ENTRIES) {
|
||||
val oldest =
|
||||
cache.entries
|
||||
.iterator()
|
||||
.next()
|
||||
.key
|
||||
cache.remove(oldest)
|
||||
}
|
||||
cache.put(LnurlForm.normalizeUrl(url), info)
|
||||
}
|
||||
|
||||
@Synchronized
|
||||
fun clear() {
|
||||
cache.clear()
|
||||
}
|
||||
|
||||
@Synchronized
|
||||
internal fun size(): Int = cache.size
|
||||
internal fun size(): Int = cache.size()
|
||||
}
|
||||
|
||||
+91
@@ -0,0 +1,91 @@
|
||||
/*
|
||||
* 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 java.util.concurrent.ConcurrentHashMap
|
||||
|
||||
/**
|
||||
* A bounded, thread-safe cache with a **lock-free [get]**.
|
||||
*
|
||||
* The common alternative — a `Collections.synchronizedMap` wrapping an
|
||||
* access-order `LinkedHashMap`, or an `@Synchronized`-guarded `LinkedHashMap` —
|
||||
* takes a single monitor on *every* operation, including reads (an access-order
|
||||
* map structurally mutates on `get`, so it cannot be read without the lock).
|
||||
* On hot read paths (zap-receipt validation, feed rich-text rendering) that one
|
||||
* monitor serializes every reader across all dispatcher threads.
|
||||
*
|
||||
* Here storage is a [ConcurrentHashMap], so [get] never takes a lock. [put] and
|
||||
* [clear] hold a small monitor that keeps the map and the recency order
|
||||
* consistent; that lock is off the read path entirely. Eviction is
|
||||
* **least-recently-put** order (a `get` does not refresh recency); re-putting a
|
||||
* key moves it to the youngest position. This matches the semantics the previous
|
||||
* monitor-based caches relied on and keeps the read path contention-free.
|
||||
*
|
||||
* Because writes and eviction happen atomically under the monitor while readers
|
||||
* see the [ConcurrentHashMap] directly, an external [size] can transiently
|
||||
* observe at most `maxSize + 1` (the instant a new entry is inserted before the
|
||||
* over-cap entry is evicted, within a single locked section). It settles to
|
||||
* `<= maxSize` once writers quiesce.
|
||||
*
|
||||
* Keys and values must be non-null (a [ConcurrentHashMap] constraint).
|
||||
*/
|
||||
class ConcurrentLruCache<K : Any, V : Any>(
|
||||
private val maxSize: Int,
|
||||
) {
|
||||
init {
|
||||
require(maxSize > 0) { "maxSize must be > 0, was $maxSize" }
|
||||
}
|
||||
|
||||
private val map = ConcurrentHashMap<K, V>()
|
||||
|
||||
// Guards writes + eviction so the map and the recency order stay consistent.
|
||||
// Reads never touch it. Writes are the cold path here, so serializing them
|
||||
// is fine; the point is a lock-free [get].
|
||||
private val writeLock = Any()
|
||||
private val order = ArrayDeque<K>()
|
||||
|
||||
fun get(key: K): V? = map[key]
|
||||
|
||||
fun put(
|
||||
key: K,
|
||||
value: V,
|
||||
) {
|
||||
synchronized(writeLock) {
|
||||
// Re-inserting an existing key makes it the youngest again.
|
||||
val existed = map.put(key, value) != null
|
||||
if (existed) order.remove(key)
|
||||
order.addLast(key)
|
||||
while (order.size > maxSize) {
|
||||
val oldest = order.removeFirst()
|
||||
map.remove(oldest)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fun clear() {
|
||||
synchronized(writeLock) {
|
||||
order.clear()
|
||||
map.clear()
|
||||
}
|
||||
}
|
||||
|
||||
fun size(): Int = map.size
|
||||
}
|
||||
Vendored
+136
@@ -0,0 +1,136 @@
|
||||
/*
|
||||
* 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 java.util.concurrent.CountDownLatch
|
||||
import java.util.concurrent.Executors
|
||||
import java.util.concurrent.TimeUnit
|
||||
import java.util.concurrent.atomic.AtomicInteger
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertNull
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
class ConcurrentLruCacheTest {
|
||||
@Test
|
||||
fun `get returns put value`() {
|
||||
val cache = ConcurrentLruCache<String, Int>(4)
|
||||
cache.put("a", 1)
|
||||
assertEquals(1, cache.get("a"))
|
||||
assertNull(cache.get("missing"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `re-put overwrites value`() {
|
||||
val cache = ConcurrentLruCache<String, Int>(4)
|
||||
cache.put("a", 1)
|
||||
cache.put("a", 2)
|
||||
assertEquals(2, cache.get("a"))
|
||||
assertEquals(1, cache.size())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `evicts least-recently-put when over capacity`() {
|
||||
val cache = ConcurrentLruCache<String, Int>(3)
|
||||
cache.put("a", 1)
|
||||
cache.put("b", 2)
|
||||
cache.put("c", 3)
|
||||
cache.put("d", 4) // pushes out "a"
|
||||
|
||||
assertNull(cache.get("a"))
|
||||
assertEquals(2, cache.get("b"))
|
||||
assertEquals(3, cache.get("c"))
|
||||
assertEquals(4, cache.get("d"))
|
||||
assertEquals(3, cache.size())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `re-putting a key refreshes its recency`() {
|
||||
val cache = ConcurrentLruCache<String, Int>(3)
|
||||
cache.put("a", 1)
|
||||
cache.put("b", 2)
|
||||
cache.put("c", 3)
|
||||
// Touch "a" via put so it is no longer the oldest.
|
||||
cache.put("a", 11)
|
||||
cache.put("d", 4) // oldest is now "b", not "a"
|
||||
|
||||
assertNull(cache.get("b"))
|
||||
assertEquals(11, cache.get("a"))
|
||||
assertEquals(3, cache.get("c"))
|
||||
assertEquals(4, cache.get("d"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `get does not refresh recency`() {
|
||||
val cache = ConcurrentLruCache<String, Int>(3)
|
||||
cache.put("a", 1)
|
||||
cache.put("b", 2)
|
||||
cache.put("c", 3)
|
||||
// A read must NOT save "a" from eviction (least-recently-put semantics).
|
||||
assertEquals(1, cache.get("a"))
|
||||
cache.put("d", 4)
|
||||
assertNull(cache.get("a"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `clear empties the cache`() {
|
||||
val cache = ConcurrentLruCache<String, Int>(4)
|
||||
cache.put("a", 1)
|
||||
cache.put("b", 2)
|
||||
cache.clear()
|
||||
assertEquals(0, cache.size())
|
||||
assertNull(cache.get("a"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `size never exceeds capacity under concurrent puts`() {
|
||||
val cap = 100
|
||||
val cache = ConcurrentLruCache<Int, Int>(cap)
|
||||
val threads = 8
|
||||
val perThread = 5_000
|
||||
val pool = Executors.newFixedThreadPool(threads)
|
||||
val start = CountDownLatch(1)
|
||||
val done = CountDownLatch(threads)
|
||||
// Writes+eviction are atomic under one lock, so an external reader can see
|
||||
// at most one over-cap entry (new inserted before old evicted).
|
||||
val overBound = AtomicInteger(0)
|
||||
|
||||
repeat(threads) { t ->
|
||||
pool.execute {
|
||||
start.await()
|
||||
for (i in 0 until perThread) {
|
||||
cache.put(t * perThread + i, i)
|
||||
// Interleave reads; they must never throw or take a lock.
|
||||
cache.get((t * perThread + i) - 1)
|
||||
if (cache.size() > cap + 1) overBound.incrementAndGet()
|
||||
}
|
||||
done.countDown()
|
||||
}
|
||||
}
|
||||
start.countDown()
|
||||
assertTrue(done.await(30, TimeUnit.SECONDS), "workers did not finish in time")
|
||||
pool.shutdown()
|
||||
|
||||
assertEquals(0, overBound.get(), "cache size exceeded cap+1 during concurrent puts")
|
||||
// Once writers quiesce it must settle to <= cap.
|
||||
assertTrue(cache.size() <= cap, "final size ${cache.size()} exceeds cap $cap")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user