mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 03:38:23 +00:00
perf(cashu): read concurrent maps through a view instead of copying them
Moving CashuWalletState off ConcurrentHashMap onto quartz ConcurrentMap turned every iteration into snapshot(), a full HashMap copy on JVM: 4-8 whole-map copies per event bundle (tokens, decrypted token contents, quotes, history, nutzaps), where the old code iterated the live views. - quartz ConcurrentMap gains asMap(): the live, weakly consistent ConcurrentHashMap on JVM/Android, and the current immutable copy-on-write map on native. Neither copies. snapshot() stays for callers that need a stable copy. ConcurrentCollectionsTest covers it. - CashuWalletState reads through asMap() everywhere it only iterates or looks up; its two session key-sets, which are copied into a new set anyway, use quartz ConcurrentSet.snapshot(). - CashuMintDirectoryState.rebuildEntries copied its announcements twice per rebuild; it now takes one snapshot, so both passes also see the same set. - The snapshot() this branch had added to the commons ConcurrentSet is reverted: quartz already has a ConcurrentSet with snapshot(), which CashuWalletState and AccountConcordActions now use. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01S7FuNBSKiyVecARSoE4B9P
This commit is contained in:
+1
-1
@@ -30,7 +30,6 @@ import com.vitorpamplona.amethyst.commons.model.cache.filter
|
||||
import com.vitorpamplona.amethyst.commons.model.concord.ConcordChannel
|
||||
import com.vitorpamplona.amethyst.commons.model.concord.ConcordCommunitySession
|
||||
import com.vitorpamplona.amethyst.commons.model.concordChannelLastReadRoute
|
||||
import com.vitorpamplona.amethyst.commons.util.ConcurrentSet
|
||||
import com.vitorpamplona.amethyst.commons.viewmodels.ReplyMode
|
||||
import com.vitorpamplona.quartz.concord.cord02Community.ConcordCommunityList.withControlRoot
|
||||
import com.vitorpamplona.quartz.concord.cord02Community.ConcordCommunityListEntry
|
||||
@@ -72,6 +71,7 @@ import com.vitorpamplona.quartz.utils.Log
|
||||
import com.vitorpamplona.quartz.utils.RandomInstance
|
||||
import com.vitorpamplona.quartz.utils.TimeUtils
|
||||
import com.vitorpamplona.quartz.utils.concurrent.ConcurrentMap
|
||||
import com.vitorpamplona.quartz.utils.concurrent.ConcurrentSet
|
||||
import kotlinx.coroutines.async
|
||||
import kotlinx.coroutines.awaitAll
|
||||
import kotlinx.coroutines.coroutineScope
|
||||
|
||||
+5
-3
@@ -217,8 +217,10 @@ class CashuMintDirectoryState(
|
||||
// ============================================================
|
||||
private fun rebuildEntries() {
|
||||
// 1. Group announcements by mint URL — most recent wins for display.
|
||||
// One copy for both passes, so byUrl and urlByDTag describe the same set of announcements.
|
||||
val announced = announcements.snapshot().values
|
||||
val byUrl: MutableMap<String, CashuMintEvent> = HashMap()
|
||||
announcements.snapshot().values.forEach { e ->
|
||||
announced.forEach { e ->
|
||||
val url = e.mintUrl() ?: return@forEach
|
||||
val existing = byUrl[url]
|
||||
if (existing == null || e.createdAt > existing.createdAt) byUrl[url] = e
|
||||
@@ -230,7 +232,7 @@ class CashuMintDirectoryState(
|
||||
// clients also store the mint URL in the `u` tag of the
|
||||
// recommendation, so we accept both.
|
||||
val urlByDTag: MutableMap<String, String> = HashMap()
|
||||
announcements.snapshot().values.forEach { e ->
|
||||
announced.forEach { e ->
|
||||
val d = e.dTag()
|
||||
val url = e.mintUrl()
|
||||
if (d != null && url != null) urlByDTag[d] = url
|
||||
@@ -241,7 +243,7 @@ class CashuMintDirectoryState(
|
||||
val perUrlFollows: MutableMap<String, MutableSet<HexKey>> = HashMap()
|
||||
val followSet = currentFollowSnapshot()
|
||||
|
||||
recommendations.snapshot().values.forEach { rec ->
|
||||
recommendations.asMap().values.forEach { rec ->
|
||||
// Mints recommended via `u` tag(s) — directly carry the URL.
|
||||
val urlsFromU = rec.mintUrls()
|
||||
// Mints recommended via `a` tag(s) — look up the URL by mint d-tag.
|
||||
|
||||
+14
-14
@@ -31,7 +31,6 @@ import com.vitorpamplona.amethyst.commons.cashu.ops.describeMintError
|
||||
import com.vitorpamplona.amethyst.commons.model.AccountSettings
|
||||
import com.vitorpamplona.amethyst.commons.model.cache.LocalCache
|
||||
import com.vitorpamplona.amethyst.commons.relayClient.assemblers.cashuProofBackfillFilters
|
||||
import com.vitorpamplona.amethyst.commons.util.ConcurrentSet
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
import com.vitorpamplona.quartz.nip01Core.core.hexToByteArray
|
||||
@@ -57,6 +56,7 @@ import com.vitorpamplona.quartz.nip61Nutzaps.nutzap.NutzapEvent
|
||||
import com.vitorpamplona.quartz.nip87Ecash.recommendation.MintRecommendationEvent
|
||||
import com.vitorpamplona.quartz.utils.Log
|
||||
import com.vitorpamplona.quartz.utils.concurrent.ConcurrentMap
|
||||
import com.vitorpamplona.quartz.utils.concurrent.ConcurrentSet
|
||||
import com.vitorpamplona.quartz.utils.secp256k1.Secp256k1
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
@@ -648,7 +648,7 @@ class CashuWalletState(
|
||||
proofBackfillDone = true
|
||||
}
|
||||
|
||||
val fresh = collected.snapshot().values.filter { tokenEvents[it.id] == null }
|
||||
val fresh = collected.asMap().values.filter { tokenEvents[it.id] == null }
|
||||
Log.i("CashuWallet") {
|
||||
"Proof backfill over ${relays.size} relay(s): ${collected.size()} kind:7375 seen, ${fresh.size} new"
|
||||
}
|
||||
@@ -779,7 +779,7 @@ class CashuWalletState(
|
||||
}
|
||||
if (dirtyTokens) recomputeUnspent()
|
||||
if (dirtyHistory) {
|
||||
_history.value = historyEvents.snapshot().values.sortedByDescending { it.createdAt }
|
||||
_history.value = historyEvents.asMap().values.sortedByDescending { it.createdAt }
|
||||
}
|
||||
if (dirtyQuotes || dirtyHistory) {
|
||||
// History gains might mark quotes as fulfilled (via the "destroyed"
|
||||
@@ -790,7 +790,7 @@ class CashuWalletState(
|
||||
triggerAutoRedeem()
|
||||
}
|
||||
if (dirtyRecommendations) {
|
||||
_ownRecommendations.value = recommendationEvents.snapshot().values.sortedByDescending { it.createdAt }
|
||||
_ownRecommendations.value = recommendationEvents.asMap().values.sortedByDescending { it.createdAt }
|
||||
}
|
||||
}
|
||||
|
||||
@@ -823,7 +823,7 @@ class CashuWalletState(
|
||||
// matching event.id and drop the entry.
|
||||
val recoKey =
|
||||
recommendationEvents
|
||||
.snapshot()
|
||||
.asMap()
|
||||
.entries
|
||||
.firstOrNull { it.value.id == id }
|
||||
?.key
|
||||
@@ -847,19 +847,19 @@ class CashuWalletState(
|
||||
settings.clearNutzapInfo()
|
||||
}
|
||||
if (dirtyTokens) recomputeUnspent()
|
||||
if (dirtyHistory) _history.value = historyEvents.snapshot().values.sortedByDescending { it.createdAt }
|
||||
if (dirtyHistory) _history.value = historyEvents.asMap().values.sortedByDescending { it.createdAt }
|
||||
if (dirtyQuotes || dirtyHistory) recomputePending()
|
||||
// dirtyNutzaps would trigger UI surfacing for inbound nutzaps; auto-
|
||||
// redeem already fires from the live-event observer, so no extra
|
||||
// signal is needed here.
|
||||
if (dirtyNutzaps) Unit
|
||||
if (dirtyRecommendations) {
|
||||
_ownRecommendations.value = recommendationEvents.snapshot().values.sortedByDescending { it.createdAt }
|
||||
_ownRecommendations.value = recommendationEvents.asMap().values.sortedByDescending { it.createdAt }
|
||||
}
|
||||
}
|
||||
|
||||
private suspend fun recomputeUnspent() {
|
||||
val all = tokenEvents.snapshot().values.toList()
|
||||
val all = tokenEvents.asMap().values.toList()
|
||||
// Decrypt anything we haven't seen before; reuse cached TokenContent
|
||||
// for events we've already decrypted. Only successes are cached, so a
|
||||
// failure is retried on the next recompute rather than being pinned as
|
||||
@@ -891,18 +891,18 @@ class CashuWalletState(
|
||||
}
|
||||
|
||||
// Shared del-rollover + sort with the headless reader.
|
||||
_tokenEntries.value = CashuWalletReader.computeUnspent(all, tokenContents.snapshot())
|
||||
_tokenEntries.value = CashuWalletReader.computeUnspent(all, tokenContents.asMap())
|
||||
}
|
||||
|
||||
/** Token events we hold but have never managed to decrypt. See [recomputeUnspent]. */
|
||||
private fun undecryptedTokenCount(): Int {
|
||||
val decrypted = tokenContents.snapshot()
|
||||
return tokenEvents.snapshot().keys.count { it !in decrypted }
|
||||
val decrypted = tokenContents.asMap()
|
||||
return tokenEvents.asMap().keys.count { it !in decrypted }
|
||||
}
|
||||
|
||||
private fun recomputePending() {
|
||||
// Shared destroyed/expired filter with the headless reader.
|
||||
_pendingQuotes.value = CashuWalletReader.computePending(quoteEvents.snapshot().values, historyEvents.snapshot().values)
|
||||
_pendingQuotes.value = CashuWalletReader.computePending(quoteEvents.asMap().values, historyEvents.asMap().values)
|
||||
}
|
||||
|
||||
private fun scanCacheForOwnEvents(): List<Event> {
|
||||
@@ -933,13 +933,13 @@ class CashuWalletState(
|
||||
// above the candidate filter needs a key, so hoist the filter.
|
||||
if (nutzapEvents.size() == 0) return
|
||||
val skipIds = HashSet<HexKey>()
|
||||
historyEvents.snapshot().values.forEach { h ->
|
||||
historyEvents.asMap().values.forEach { h ->
|
||||
h.redeemedReferences().forEach { skipIds.add(it.eventId) }
|
||||
}
|
||||
skipIds.addAll(sessionRedeemedNutzaps.snapshot())
|
||||
skipIds.addAll(sessionUnredeemableNutzaps.snapshot())
|
||||
|
||||
val candidates = nutzapEvents.snapshot().values.filter { it.id !in skipIds }
|
||||
val candidates = nutzapEvents.asMap().values.filter { it.id !in skipIds }
|
||||
if (candidates.isEmpty()) return
|
||||
|
||||
val privkey = walletPrivkeyHex() ?: return
|
||||
|
||||
@@ -50,7 +50,4 @@ expect class ConcurrentSet<E : Any>() {
|
||||
fun clear()
|
||||
|
||||
val size: Int
|
||||
|
||||
/** A point-in-time copy of the elements — safe to iterate while others add and remove. */
|
||||
fun snapshot(): Set<E>
|
||||
}
|
||||
|
||||
-12
@@ -60,16 +60,4 @@ class ConcurrentSetTest {
|
||||
assertEquals(0, set.size)
|
||||
assertFalse(set.contains("a"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `snapshot is a copy that later writes do not change`() {
|
||||
val set = ConcurrentSet<String>()
|
||||
set.add("a")
|
||||
set.add("b")
|
||||
val snapshot = set.snapshot()
|
||||
set.add("c")
|
||||
set.remove("a")
|
||||
assertEquals(setOf("a", "b"), snapshot)
|
||||
assertEquals(setOf("b", "c"), set.snapshot())
|
||||
}
|
||||
}
|
||||
|
||||
-2
@@ -36,6 +36,4 @@ actual class ConcurrentSet<E : Any> {
|
||||
actual fun clear() = lock.withLock { set.clear() }
|
||||
|
||||
actual val size: Int get() = lock.withLock { set.size }
|
||||
|
||||
actual fun snapshot(): Set<E> = lock.withLock { set.toHashSet() }
|
||||
}
|
||||
|
||||
-3
@@ -35,7 +35,4 @@ actual class ConcurrentSet<E : Any> {
|
||||
actual fun clear() = set.clear()
|
||||
|
||||
actual val size: Int get() = set.size
|
||||
|
||||
// The key-set view iterates weakly consistently, so copying it never throws while others write.
|
||||
actual fun snapshot(): Set<E> = set.toHashSet()
|
||||
}
|
||||
|
||||
+9
@@ -93,4 +93,13 @@ expect class ConcurrentMap<K : Any, V : Any>() {
|
||||
|
||||
/** A point-in-time copy of the entries — safe to iterate without holding a lock. */
|
||||
fun snapshot(): Map<K, V>
|
||||
|
||||
/**
|
||||
* A read-only view to iterate or look up in right away, without copying: the live,
|
||||
* weakly consistent map on JVM/Android (iteration never throws while others write, and may
|
||||
* or may not see those writes), the current immutable copy-on-write state on native. Don't
|
||||
* hold it expecting later writes to show up; use [snapshot] when a caller needs a stable
|
||||
* copy it can keep.
|
||||
*/
|
||||
fun asMap(): Map<K, V>
|
||||
}
|
||||
|
||||
+14
@@ -118,6 +118,20 @@ class ConcurrentCollectionsTest {
|
||||
assertEquals(3, m.size())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun mapAsMapSeesEarlierWritesAndSurvivesLaterOnes() {
|
||||
val map = ConcurrentMap<String, Int>()
|
||||
map["a"] = 1
|
||||
map["b"] = 2
|
||||
val view = map.asMap()
|
||||
assertEquals(mapOf("a" to 1, "b" to 2), view.toMap())
|
||||
// Writing while iterating must not throw on any target. One fixed key, so a view that
|
||||
// does see the new entry (JVM may) still ends.
|
||||
for (key in view.keys) map["written-during-iteration"] = 0
|
||||
assertEquals(0, map["written-during-iteration"])
|
||||
assertEquals(1, map["a"])
|
||||
}
|
||||
|
||||
@Test
|
||||
fun setAddContainsSize() {
|
||||
val s = ConcurrentSet<String>()
|
||||
|
||||
+2
@@ -66,4 +66,6 @@ actual class ConcurrentMap<K : Any, V : Any> {
|
||||
actual fun size(): Int = map.size
|
||||
|
||||
actual fun snapshot(): Map<K, V> = HashMap(map)
|
||||
|
||||
actual fun asMap(): Map<K, V> = map
|
||||
}
|
||||
|
||||
+3
@@ -119,4 +119,7 @@ actual class ConcurrentMap<K : Any, V : Any> {
|
||||
actual fun size(): Int = ref.load().size
|
||||
|
||||
actual fun snapshot(): Map<K, V> = HashMap(ref.load())
|
||||
|
||||
// The published map is never mutated after the CAS that installed it, so it is its own view.
|
||||
actual fun asMap(): Map<K, V> = ref.load()
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user