mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
Merge pull request #4313 from vitorpamplona/claude/adoring-planck-mewnwe
Unify Apple and Linux LargeCache to shared Kotlin/Native impl
This commit is contained in:
+1
-1
@@ -70,7 +70,7 @@ class ChessEventCollector(
|
|||||||
val startEvent: StateFlow<JesterEvent?> = _startEvent.asStateFlow()
|
val startEvent: StateFlow<JesterEvent?> = _startEvent.asStateFlow()
|
||||||
|
|
||||||
// Move events (deduplicated by event ID). String keys are Comparable, so
|
// Move events (deduplicated by event ID). String keys are Comparable, so
|
||||||
// LargeCache (ConcurrentSkipListMap on JVM, CacheMap on Apple) works.
|
// LargeCache (ConcurrentSkipListMap on JVM, StripedHashMap on native) works.
|
||||||
private val moves = LargeCache<String, JesterEvent>()
|
private val moves = LargeCache<String, JesterEvent>()
|
||||||
|
|
||||||
// Track all processed event IDs for fast deduplication. A plain HashSet
|
// Track all processed event IDs for fast deduplication. A plain HashSet
|
||||||
|
|||||||
@@ -10,7 +10,6 @@ accompanistAdaptive = "0.37.3"
|
|||||||
# LargeCache/ConcurrentHashCache — they read those views directly and through the
|
# LargeCache/ConcurrentHashCache — they read those views directly and through the
|
||||||
# stdlib Map operators (filter/map/groupBy/associate/count). Bumping needs a port to
|
# stdlib Map operators (filter/map/groupBy/associate/count). Bumping needs a port to
|
||||||
# the new forEach/forEachKey/forEachValue API, verified on an Apple target.
|
# the new forEach/forEachKey/forEachValue API, verified on an Apple target.
|
||||||
cachemapVersion = "0.2.4"
|
|
||||||
composeMultiplatform = "1.12.1"
|
composeMultiplatform = "1.12.1"
|
||||||
activityCompose = "1.13.0"
|
activityCompose = "1.13.0"
|
||||||
# 9.4.0's PerModuleBundleTask rejects an AAB entry whose name contains a colon, which every
|
# 9.4.0's PerModuleBundleTask rejects an AAB entry whose name contains a colon, which every
|
||||||
@@ -174,7 +173,6 @@ androidx-ui-test-junit4 = { group = "androidx.compose.ui", name = "ui-test-junit
|
|||||||
androidx-ui-tooling = { group = "androidx.compose.ui", name = "ui-tooling" }
|
androidx-ui-tooling = { group = "androidx.compose.ui", name = "ui-tooling" }
|
||||||
androidx-ui-tooling-preview = { group = "androidx.compose.ui", name = "ui-tooling-preview" }
|
androidx-ui-tooling-preview = { group = "androidx.compose.ui", name = "ui-tooling-preview" }
|
||||||
audiowaveform = { group = "com.github.lincollincol", name = "compose-audiowaveform", version.ref = "audiowaveform" }
|
audiowaveform = { group = "com.github.lincollincol", name = "compose-audiowaveform", version.ref = "audiowaveform" }
|
||||||
charlietap-cachemap = { module = "io.github.charlietap:cachemap", version.ref = "cachemapVersion" }
|
|
||||||
coil-compose = { group = "io.coil-kt.coil3", name = "coil-compose", version.ref = "coil" }
|
coil-compose = { group = "io.coil-kt.coil3", name = "coil-compose", version.ref = "coil" }
|
||||||
coil-gif = { group = "io.coil-kt.coil3", name = "coil-gif", version.ref = "coil" }
|
coil-gif = { group = "io.coil-kt.coil3", name = "coil-gif", version.ref = "coil" }
|
||||||
coil-svg = { group = "io.coil-kt.coil3", name = "coil-svg", version.ref = "coil" }
|
coil-svg = { group = "io.coil-kt.coil3", name = "coil-svg", version.ref = "coil" }
|
||||||
|
|||||||
@@ -296,7 +296,6 @@ kotlin {
|
|||||||
create("appleMain") {
|
create("appleMain") {
|
||||||
dependsOn(nativeMain)
|
dependsOn(nativeMain)
|
||||||
dependencies {
|
dependencies {
|
||||||
implementation(libs.charlietap.cachemap)
|
|
||||||
implementation(libs.dev.whyoleg.cryptography.provider.apple.optimal)
|
implementation(libs.dev.whyoleg.cryptography.provider.apple.optimal)
|
||||||
implementation("io.github.andreypfau:kotlinx-crypto-hmac:0.0.4")
|
implementation("io.github.andreypfau:kotlinx-crypto-hmac:0.0.4")
|
||||||
implementation("io.github.andreypfau:kotlinx-crypto-sha2:0.0.4")
|
implementation("io.github.andreypfau:kotlinx-crypto-sha2:0.0.4")
|
||||||
|
|||||||
Vendored
-40
@@ -1,40 +0,0 @@
|
|||||||
/*
|
|
||||||
* 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 io.github.charlietap.cachemap.cacheMapOf
|
|
||||||
|
|
||||||
actual class ConcurrentHashCache<K : Any, V : Any> {
|
|
||||||
private val map = cacheMapOf<K, V>()
|
|
||||||
|
|
||||||
actual fun get(key: K): V? = map[key]
|
|
||||||
|
|
||||||
actual fun put(
|
|
||||||
key: K,
|
|
||||||
value: V,
|
|
||||||
) {
|
|
||||||
map.put(key, value)
|
|
||||||
}
|
|
||||||
|
|
||||||
actual fun size(): Int = map.size
|
|
||||||
|
|
||||||
actual fun clear() = map.clear()
|
|
||||||
}
|
|
||||||
-432
@@ -1,432 +0,0 @@
|
|||||||
/*
|
|
||||||
* 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 io.github.charlietap.cachemap.CacheMap
|
|
||||||
import io.github.charlietap.cachemap.cacheMapOf
|
|
||||||
|
|
||||||
// An implementation of a Threadsafe map, using CacheMap.
|
|
||||||
// Investigating a Swift-based alternative(for now)
|
|
||||||
actual class LargeCache<K, V> : ICacheOperations<K, V> {
|
|
||||||
private val concurrentMap = cacheMapOf<K, V>()
|
|
||||||
|
|
||||||
actual fun keys(): Set<K> = concurrentMap.keys
|
|
||||||
|
|
||||||
actual fun values(): Iterable<V> = concurrentMap.values
|
|
||||||
|
|
||||||
actual fun get(key: K): V? = concurrentMap[key]
|
|
||||||
|
|
||||||
actual fun remove(key: K): V? = concurrentMap.remove(key)
|
|
||||||
|
|
||||||
actual fun isEmpty(): Boolean = concurrentMap.isEmpty()
|
|
||||||
|
|
||||||
actual fun clear() {
|
|
||||||
concurrentMap.clear()
|
|
||||||
}
|
|
||||||
|
|
||||||
actual fun containsKey(key: K): Boolean = concurrentMap.containsKey(key)
|
|
||||||
|
|
||||||
actual fun put(
|
|
||||||
key: K,
|
|
||||||
value: V,
|
|
||||||
) {
|
|
||||||
concurrentMap.put(key, value)
|
|
||||||
}
|
|
||||||
|
|
||||||
actual fun getOrCreate(
|
|
||||||
key: K,
|
|
||||||
builder: (key: K) -> V,
|
|
||||||
): V {
|
|
||||||
val value = concurrentMap.get(key)
|
|
||||||
|
|
||||||
return if (value != null) {
|
|
||||||
value
|
|
||||||
} else {
|
|
||||||
val newObject = builder(key)
|
|
||||||
concurrentMap.put(key, newObject)
|
|
||||||
concurrentMap[key] ?: newObject
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
actual fun createIfAbsent(
|
|
||||||
key: K,
|
|
||||||
builder: (key: K) -> V,
|
|
||||||
): Boolean {
|
|
||||||
val value = concurrentMap.get(key)
|
|
||||||
return if (value != null) {
|
|
||||||
false
|
|
||||||
} else {
|
|
||||||
val newObject = builder(key)
|
|
||||||
concurrentMap.put(key, newObject)
|
|
||||||
concurrentMap[key] != null
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
actual override fun size(): Int = concurrentMap.size
|
|
||||||
|
|
||||||
actual override fun forEach(consumer: ICacheBiConsumer<K, V>) {
|
|
||||||
// Take a snapshot of entries to avoid ConcurrentModificationException
|
|
||||||
// when the map is modified during iteration (e.g., NostrClient.syncFilters
|
|
||||||
// iterating while subscriptions are added from another coroutine).
|
|
||||||
concurrentMap.entries.toList().forEach { consumer.accept(it.key, it.value) }
|
|
||||||
}
|
|
||||||
|
|
||||||
actual override fun filter(consumer: CacheCollectors.BiFilter<K, V>): List<V> =
|
|
||||||
concurrentMap
|
|
||||||
.filter { consumer.filter(it.key, it.value) }
|
|
||||||
.values
|
|
||||||
.toList()
|
|
||||||
|
|
||||||
actual override fun filterIntoSet(consumer: CacheCollectors.BiFilter<K, V>): Set<V> =
|
|
||||||
concurrentMap
|
|
||||||
.filter { consumer.filter(it.key, it.value) }
|
|
||||||
.values
|
|
||||||
.toSet()
|
|
||||||
|
|
||||||
actual override fun <R> map(consumer: CacheCollectors.BiNotNullMapper<K, V, R>): List<R> = concurrentMap.map { consumer.map(it.key, it.value) }
|
|
||||||
|
|
||||||
actual override fun <R> mapNotNull(consumer: CacheCollectors.BiMapper<K, V, R?>): List<R> = concurrentMap.mapNotNull { consumer.map(it.key, it.value) }
|
|
||||||
|
|
||||||
actual override fun <R> mapNotNullIntoSet(consumer: CacheCollectors.BiMapper<K, V, R?>): Set<R> = mapNotNull(consumer).toSet()
|
|
||||||
|
|
||||||
actual override fun <R> mapFlatten(consumer: CacheCollectors.BiMapper<K, V, Collection<R>?>): List<R> = concurrentMap.flatMap { entry -> consumer.map(entry.key, entry.value) ?: emptyList() }
|
|
||||||
|
|
||||||
actual override fun <R> mapFlattenIntoSet(consumer: CacheCollectors.BiMapper<K, V, Collection<R>?>): Set<R> = mapFlatten(consumer).toSet()
|
|
||||||
|
|
||||||
actual override fun maxOrNullOf(
|
|
||||||
filter: CacheCollectors.BiFilter<K, V>,
|
|
||||||
comparator: Comparator<V>,
|
|
||||||
): V? {
|
|
||||||
// return concurrentMap.maxOfWithOrNull(
|
|
||||||
// comparator,
|
|
||||||
// selector = {
|
|
||||||
// if (filter.filter(it.key, it.value)) it.value else concurrentMap.getValue(it.key)
|
|
||||||
// }
|
|
||||||
// )
|
|
||||||
return concurrentMap.maxOrNullOf(filter, comparator)
|
|
||||||
}
|
|
||||||
|
|
||||||
actual override fun sumOf(consumer: CacheCollectors.BiSumOf<K, V>): Int {
|
|
||||||
return concurrentMap.sumOf(consumer)
|
|
||||||
// return concurrentMap.map { consumer.map(it.key, it.value) }.sum()
|
|
||||||
}
|
|
||||||
|
|
||||||
actual override fun sumOfLong(consumer: CacheCollectors.BiSumOfLong<K, V>): Long = concurrentMap.sumOfLong(consumer)
|
|
||||||
|
|
||||||
actual override fun <R> groupBy(consumer: CacheCollectors.BiNotNullMapper<K, V, R>): Map<R, List<V>> = concurrentMap.groupBy(consumer)
|
|
||||||
|
|
||||||
actual override fun <R> countByGroup(consumer: CacheCollectors.BiNotNullMapper<K, V, R>): Map<R, Int> = concurrentMap.countByGroup(consumer)
|
|
||||||
|
|
||||||
actual override fun <R> sumByGroup(
|
|
||||||
groupMap: CacheCollectors.BiNotNullMapper<K, V, R>,
|
|
||||||
sumOf: CacheCollectors.BiNotNullMapper<K, V, Long>,
|
|
||||||
): Map<R, Long> = concurrentMap.sumByGroup(groupMap, sumOf)
|
|
||||||
|
|
||||||
actual override fun count(consumer: CacheCollectors.BiFilter<K, V>): Int = concurrentMap.count { consumer.filter(it.key, it.value) }
|
|
||||||
|
|
||||||
actual override fun <T, U> associate(transform: (K, V) -> Pair<T, U>): Map<T, U> = concurrentMap.associate(transform)
|
|
||||||
|
|
||||||
actual override fun <U> associateWith(transform: (K, V) -> U?): Map<K, U?> = concurrentMap.associateWith(transform)
|
|
||||||
|
|
||||||
actual override fun filter(
|
|
||||||
from: K,
|
|
||||||
to: K,
|
|
||||||
consumer: CacheCollectors.BiFilter<K, V>,
|
|
||||||
): List<V> {
|
|
||||||
val transientList = concurrentMap.subMapAlt(from, to)
|
|
||||||
return transientList.filter { consumer.filter(it.key, it.value) }.values.toList()
|
|
||||||
}
|
|
||||||
|
|
||||||
actual override fun filterIntoSet(
|
|
||||||
from: K,
|
|
||||||
to: K,
|
|
||||||
consumer: CacheCollectors.BiFilter<K, V>,
|
|
||||||
): Set<V> = filter(from, to, consumer).toSet()
|
|
||||||
|
|
||||||
actual override fun <R> map(
|
|
||||||
from: K,
|
|
||||||
to: K,
|
|
||||||
consumer: CacheCollectors.BiNotNullMapper<K, V, R>,
|
|
||||||
): List<R> {
|
|
||||||
val transientList = concurrentMap.subMapAlt(from, to)
|
|
||||||
return transientList.map { consumer.map(it.key, it.value) }
|
|
||||||
}
|
|
||||||
|
|
||||||
actual override fun <R> mapNotNull(
|
|
||||||
from: K,
|
|
||||||
to: K,
|
|
||||||
consumer: CacheCollectors.BiMapper<K, V, R?>,
|
|
||||||
): List<R> = concurrentMap.subMapAlt(from, to).mapNotNull { consumer.map(it.key, it.value) }
|
|
||||||
|
|
||||||
actual override fun <R> mapNotNullIntoSet(
|
|
||||||
from: K,
|
|
||||||
to: K,
|
|
||||||
consumer: CacheCollectors.BiMapper<K, V, R?>,
|
|
||||||
): Set<R> = mapNotNull(from, to, consumer).toSet()
|
|
||||||
|
|
||||||
actual override fun <R> mapFlatten(
|
|
||||||
from: K,
|
|
||||||
to: K,
|
|
||||||
consumer: CacheCollectors.BiMapper<K, V, Collection<R>?>,
|
|
||||||
): List<R> = concurrentMap.subMapAlt(from, to).flatMap { consumer.map(it.key, it.value) as Iterable<R> }
|
|
||||||
|
|
||||||
actual override fun <R> mapFlattenIntoSet(
|
|
||||||
from: K,
|
|
||||||
to: K,
|
|
||||||
consumer: CacheCollectors.BiMapper<K, V, Collection<R>?>,
|
|
||||||
): Set<R> = mapFlatten(from, to, consumer).toSet()
|
|
||||||
|
|
||||||
actual override fun maxOrNullOf(
|
|
||||||
from: K,
|
|
||||||
to: K,
|
|
||||||
filter: CacheCollectors.BiFilter<K, V>,
|
|
||||||
comparator: Comparator<V>,
|
|
||||||
): V? {
|
|
||||||
val transient = concurrentMap.subMapAlt(from, to)
|
|
||||||
return transient.maxOrNullOf(filter, comparator)
|
|
||||||
}
|
|
||||||
|
|
||||||
actual override fun sumOf(
|
|
||||||
from: K,
|
|
||||||
to: K,
|
|
||||||
consumer: CacheCollectors.BiSumOf<K, V>,
|
|
||||||
): Int = concurrentMap.subMapAlt(from, to).sumOf(consumer)
|
|
||||||
|
|
||||||
actual override fun sumOfLong(
|
|
||||||
from: K,
|
|
||||||
to: K,
|
|
||||||
consumer: CacheCollectors.BiSumOfLong<K, V>,
|
|
||||||
): Long = concurrentMap.subMapAlt(from, to).sumOfLong(consumer)
|
|
||||||
|
|
||||||
actual override fun <R> groupBy(
|
|
||||||
from: K,
|
|
||||||
to: K,
|
|
||||||
consumer: CacheCollectors.BiNotNullMapper<K, V, R>,
|
|
||||||
): Map<R, List<V>> = concurrentMap.subMapAlt(from, to).groupBy(consumer)
|
|
||||||
|
|
||||||
actual override fun <R> countByGroup(
|
|
||||||
from: K,
|
|
||||||
to: K,
|
|
||||||
consumer: CacheCollectors.BiNotNullMapper<K, V, R>,
|
|
||||||
): Map<R, Int> = concurrentMap.subMapAlt(from, to).countByGroup(consumer)
|
|
||||||
|
|
||||||
actual override fun <R> sumByGroup(
|
|
||||||
from: K,
|
|
||||||
to: K,
|
|
||||||
groupMap: CacheCollectors.BiNotNullMapper<K, V, R>,
|
|
||||||
sumOf: CacheCollectors.BiNotNullMapper<K, V, Long>,
|
|
||||||
): Map<R, Long> = concurrentMap.subMapAlt(from, to).sumByGroup(groupMap, sumOf)
|
|
||||||
|
|
||||||
actual override fun count(
|
|
||||||
from: K,
|
|
||||||
to: K,
|
|
||||||
consumer: CacheCollectors.BiFilter<K, V>,
|
|
||||||
): Int = concurrentMap.subMapAlt(from, to).count { consumer.filter(it.key, it.value) }
|
|
||||||
|
|
||||||
actual override fun <T, U> associate(
|
|
||||||
from: K,
|
|
||||||
to: K,
|
|
||||||
transform: (K, V) -> Pair<T, U>,
|
|
||||||
): Map<T, U> = concurrentMap.subMapAlt(from, to).associate(transform)
|
|
||||||
|
|
||||||
actual override fun <U> associateWith(
|
|
||||||
from: K,
|
|
||||||
to: K,
|
|
||||||
transform: (K, V) -> U?,
|
|
||||||
): Map<K, U?> = concurrentMap.subMapAlt(from, to).associateWith(transform)
|
|
||||||
|
|
||||||
actual override fun joinToString(
|
|
||||||
separator: CharSequence,
|
|
||||||
prefix: CharSequence,
|
|
||||||
postfix: CharSequence,
|
|
||||||
limit: Int,
|
|
||||||
truncated: CharSequence,
|
|
||||||
transform: ((K, V) -> CharSequence)?,
|
|
||||||
): String {
|
|
||||||
val buffer = StringBuilder()
|
|
||||||
buffer.append(prefix)
|
|
||||||
var count = 0
|
|
||||||
forEach { key, value ->
|
|
||||||
val str = if (transform != null) transform(key, value) else ""
|
|
||||||
if (str.isNotEmpty()) {
|
|
||||||
if (++count > 1) buffer.append(separator)
|
|
||||||
if (limit < 0 || count <= limit) {
|
|
||||||
when {
|
|
||||||
transform != null -> buffer.append(str)
|
|
||||||
else -> buffer.append("$key $value")
|
|
||||||
}
|
|
||||||
} else {
|
|
||||||
return@forEach
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if (limit >= 0 && count > limit) buffer.append(truncated)
|
|
||||||
buffer.append(postfix)
|
|
||||||
return buffer.toString()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Different subMap implementations below. Investigating their performance for now.
|
|
||||||
|
|
||||||
fun <K, V> CacheMap<K, V>.subMapSlow(
|
|
||||||
from: K,
|
|
||||||
to: K,
|
|
||||||
toInclusive: Boolean = true,
|
|
||||||
): Map<K, V> {
|
|
||||||
val transientList = toList()
|
|
||||||
val transientSubList =
|
|
||||||
transientList.subList(
|
|
||||||
fromIndex = transientList.indexOf(Pair(from, getValue(from))),
|
|
||||||
toIndex = transientList.indexOf(Pair(to, getValue(to))),
|
|
||||||
)
|
|
||||||
val completeSubList = transientSubList + Pair(to, getValue(to))
|
|
||||||
|
|
||||||
return if (toInclusive) completeSubList.toMap() else transientSubList.toMap()
|
|
||||||
}
|
|
||||||
|
|
||||||
fun <K, V> CacheMap<K, V>.subMapAlt(
|
|
||||||
from: K,
|
|
||||||
to: K,
|
|
||||||
toInclusive: Boolean = true,
|
|
||||||
): Map<K, V> {
|
|
||||||
val resultMap = hashMapOf<K, V>()
|
|
||||||
val keySet = keys
|
|
||||||
val fromIndex = keySet.indexOf(from)
|
|
||||||
val toIndex = keySet.indexOf(to)
|
|
||||||
for (index in fromIndex until toIndex) {
|
|
||||||
val correspondingEntry = entries.elementAt(index)
|
|
||||||
resultMap[correspondingEntry.key] = correspondingEntry.value
|
|
||||||
}
|
|
||||||
if (toInclusive) {
|
|
||||||
val correspondingToEntry = entries.elementAt(toIndex)
|
|
||||||
resultMap[correspondingToEntry.key] = correspondingToEntry.value
|
|
||||||
}
|
|
||||||
|
|
||||||
return resultMap
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* The following functions below are (re)implementations for the ICacheOperations
|
|
||||||
* interface. A lot of it is copying and pasting, with modifications to make it work
|
|
||||||
* consistently.
|
|
||||||
*/
|
|
||||||
|
|
||||||
fun <K, V> Map<K, V>.maxOrNullOf(
|
|
||||||
filter: CacheCollectors.BiFilter<K, V>,
|
|
||||||
comparator: Comparator<V>,
|
|
||||||
): V? {
|
|
||||||
var maxK: K? = null
|
|
||||||
var maxV: V? = null
|
|
||||||
forEach {
|
|
||||||
if (filter.filter(it.key, it.value)) {
|
|
||||||
if (maxK == null || (maxV != null && comparator.compare(it.value, maxV) > 0)) {
|
|
||||||
maxK = it.key
|
|
||||||
maxV = it.value
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
val finalMaxK: K? = maxK
|
|
||||||
val finalMaxV: V? = maxV
|
|
||||||
|
|
||||||
return finalMaxV
|
|
||||||
}
|
|
||||||
|
|
||||||
fun <K, V> Map<K, V>.sumOf(consumer: CacheCollectors.BiSumOf<K, V>): Int {
|
|
||||||
var sum = 0
|
|
||||||
forEach { sum += consumer.map(it.key, it.value) }
|
|
||||||
return sum
|
|
||||||
}
|
|
||||||
|
|
||||||
fun <K, V> Map<K, V>.sumOfLong(consumer: CacheCollectors.BiSumOfLong<K, V>): Long {
|
|
||||||
var sum = 0L
|
|
||||||
forEach { sum += consumer.map(it.key, it.value) }
|
|
||||||
return sum
|
|
||||||
}
|
|
||||||
|
|
||||||
fun <K, V, R> Map<K, V>.groupBy(consumer: CacheCollectors.BiNotNullMapper<K, V, R>): Map<R, List<V>> {
|
|
||||||
val results = HashMap<R, ArrayList<V>>()
|
|
||||||
forEach {
|
|
||||||
val group = consumer.map(it.key, it.value)
|
|
||||||
val list = results[group]
|
|
||||||
if (list == null) {
|
|
||||||
val answer = ArrayList<V>()
|
|
||||||
answer.add(it.value)
|
|
||||||
results[group] = answer
|
|
||||||
} else {
|
|
||||||
list.add(it.value)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
return results
|
|
||||||
}
|
|
||||||
|
|
||||||
fun <K, V, R> Map<K, V>.countByGroup(consumer: CacheCollectors.BiNotNullMapper<K, V, R>): Map<R, Int> {
|
|
||||||
val results = HashMap<R, Int>()
|
|
||||||
forEach {
|
|
||||||
val group = consumer.map(it.key, it.value)
|
|
||||||
val count = results[group]
|
|
||||||
if (count == null) {
|
|
||||||
results[group] = 1
|
|
||||||
} else {
|
|
||||||
results[group] = count + 1
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
return results
|
|
||||||
}
|
|
||||||
|
|
||||||
fun <K, V, R> Map<K, V>.sumByGroup(
|
|
||||||
groupMap: CacheCollectors.BiNotNullMapper<K, V, R>,
|
|
||||||
sumOf: CacheCollectors.BiNotNullMapper<K, V, Long>,
|
|
||||||
): Map<R, Long> {
|
|
||||||
val results = HashMap<R, Long>()
|
|
||||||
forEach {
|
|
||||||
val group = groupMap.map(it.key, it.value)
|
|
||||||
val sum = results[group]
|
|
||||||
if (sum == null) {
|
|
||||||
results[group] = sumOf.map(it.key, it.value)
|
|
||||||
} else {
|
|
||||||
results[group] = sum + sumOf.map(it.key, it.value)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
return results
|
|
||||||
}
|
|
||||||
|
|
||||||
fun <K, V, T, U> Map<K, V>.associate(transform: (K, V) -> Pair<T, U>): Map<T, U> {
|
|
||||||
val results: LinkedHashMap<T, U> = LinkedHashMap(size)
|
|
||||||
forEach {
|
|
||||||
val pair = transform(it.key, it.value)
|
|
||||||
results[pair.first] = pair.second
|
|
||||||
}
|
|
||||||
|
|
||||||
return results
|
|
||||||
}
|
|
||||||
|
|
||||||
fun <K, V, U> Map<K, V>.associateWith(transform: (K, V) -> U?): Map<K, U?> {
|
|
||||||
val results: LinkedHashMap<K, U?> = LinkedHashMap(size)
|
|
||||||
forEach {
|
|
||||||
results[it.key] = transform(it.key, it.value)
|
|
||||||
}
|
|
||||||
|
|
||||||
return results
|
|
||||||
}
|
|
||||||
+1
-1
@@ -106,7 +106,7 @@ class RelayAuthenticator(
|
|||||||
) : IAuthStatus {
|
) : IAuthStatus {
|
||||||
// Connection callbacks fire on the per-relay OkHttp dispatcher thread, so
|
// Connection callbacks fire on the per-relay OkHttp dispatcher thread, so
|
||||||
// this state is mutated concurrently — LargeCache wraps a platform-tuned
|
// this state is mutated concurrently — LargeCache wraps a platform-tuned
|
||||||
// concurrent map (ConcurrentSkipListMap on jvmAndroid, CacheMap on Apple).
|
// concurrent map (ConcurrentSkipListMap on jvmAndroid, StripedHashMap on native).
|
||||||
//
|
//
|
||||||
// This stays mutable because RelayAuthStatus carries an LruCache that has
|
// This stays mutable because RelayAuthStatus carries an LruCache that has
|
||||||
// to be addressable from the dispatcher thread. The Compose-observable
|
// to be addressable from the dispatcher thread. The Compose-observable
|
||||||
|
|||||||
+2
-2
@@ -30,8 +30,8 @@ package com.vitorpamplona.quartz.utils.cache
|
|||||||
* duplicate check in
|
* duplicate check in
|
||||||
* [com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.CachingEventDecoder].
|
* [com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.CachingEventDecoder].
|
||||||
*
|
*
|
||||||
* Actuals: JVM/Android → `ConcurrentHashMap`; Apple → `CacheMap` (same
|
* Actuals: JVM/Android → `ConcurrentHashMap`; Apple and Linux → a striped-lock
|
||||||
* backing as [LargeCache]); Linux → copy-on-write (CI-only target).
|
* chained hash table (`StripedHashMap`, same backing as [LargeCache]).
|
||||||
*/
|
*/
|
||||||
expect class ConcurrentHashCache<K : Any, V : Any>() {
|
expect class ConcurrentHashCache<K : Any, V : Any>() {
|
||||||
fun get(key: K): V?
|
fun get(key: K): V?
|
||||||
|
|||||||
+6
-8
@@ -62,20 +62,18 @@ import kotlin.time.TimeSource
|
|||||||
* two HashMaps per put, 60k live ratio 2.28 - 5.21
|
* two HashMaps per put, 60k live ratio 2.28 - 5.21
|
||||||
* ```
|
* ```
|
||||||
*
|
*
|
||||||
* The last row is what Apple targets actually run: LargeCache there wraps
|
* The last row is what Apple targets used to run: LargeCache there wrapped
|
||||||
* charlietap's CacheMap, whose LeftRight `mutate` applies each write to both
|
* charlietap's CacheMap, whose LeftRight `mutate` applies each write to both
|
||||||
* of its two maps under a lock — O(1), no copying. It still crossed the 5.0
|
* of its two maps under a lock — O(1), no copying. It still crossed the 5.0
|
||||||
* threshold on a loaded machine, which is how this failed on the iOS simulator
|
* threshold on a loaded machine, which is how this failed on the iOS simulator
|
||||||
* without anything being wrong with the outbox. The `retains nothing` row is
|
* without anything being wrong with the outbox. The `retains nothing` row is
|
||||||
* the control: same allocations, nothing kept alive, ratio flat.
|
* the control: same allocations, nothing kept alive, ratio flat.
|
||||||
*
|
*
|
||||||
* So on Apple the O(1) guarantee comes from the data structure by
|
* Kotlin/Native (Apple and Linux) now backs LargeCache with StripedHashMap,
|
||||||
* construction and does not need pinning here; on JVM/Android it comes from
|
* O(1) per write by construction, so the guarantee does not need pinning
|
||||||
* LargeCache being a ConcurrentHashMap, and a regression in this class
|
* there; on JVM/Android it comes from LargeCache being a concurrent map, and a
|
||||||
* (someone reintroducing a copy-on-write map, or a per-publish full scan)
|
* regression in this class (someone reintroducing a copy-on-write map, or a
|
||||||
* shows up cleanly. Note that Kotlin/Native's linuxX64 LargeCache IS
|
* per-publish full scan) shows up cleanly.
|
||||||
* copy-on-write today — that is a real cost, but not one a wall-clock ratio
|
|
||||||
* can report reliably, as the numbers above show.
|
|
||||||
*/
|
*/
|
||||||
class PoolEventOutboxScaleTest {
|
class PoolEventOutboxScaleTest {
|
||||||
private val relay = NormalizedRelayUrl("wss://scale.relay.test")
|
private val relay = NormalizedRelayUrl("wss://scale.relay.test")
|
||||||
|
|||||||
+2
-2
@@ -21,8 +21,8 @@
|
|||||||
package com.vitorpamplona.quartz.utils.cache
|
package com.vitorpamplona.quartz.utils.cache
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Linux/Native actual for [ConcurrentHashCache], over the same [StripedHashMap] as
|
* Kotlin/Native actual (Apple and Linux) for [ConcurrentHashCache], over the same
|
||||||
* `LargeCache.linux.kt` — read its docs for why.
|
* [StripedHashMap] as `LargeCache.native.kt` — read its docs for why.
|
||||||
*
|
*
|
||||||
* This one was the worst-placed of the copy-on-write caches: its only caller,
|
* 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
|
* `CachingEventDecoder`, writes once per event arriving from a relay, so every decode
|
||||||
+11
-2
@@ -21,7 +21,8 @@
|
|||||||
package com.vitorpamplona.quartz.utils.cache
|
package com.vitorpamplona.quartz.utils.cache
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Linux/Native actual for [LargeCache] — the store behind Amethyst's `LocalCache`.
|
* Kotlin/Native actual (Apple and Linux) for [LargeCache] — the store behind Amethyst's
|
||||||
|
* `LocalCache`.
|
||||||
*
|
*
|
||||||
* All of the concurrency and the performance rationale lives in [StripedHashMap]; this
|
* All of the concurrency and the performance rationale lives in [StripedHashMap]; this
|
||||||
* class is only the [ICacheOperations] surface over it. The short version: `LocalCache`
|
* class is only the [ICacheOperations] surface over it. The short version: `LocalCache`
|
||||||
@@ -40,8 +41,16 @@ package com.vitorpamplona.quartz.utils.cache
|
|||||||
* - There is no snapshot and no defensive `entries.toList()`, so no
|
* - There is no snapshot and no defensive `entries.toList()`, so no
|
||||||
* `ConcurrentModificationException` window and no per-scan copy.
|
* `ConcurrentModificationException` window and no per-scan copy.
|
||||||
*
|
*
|
||||||
|
* Apple used to wrap charlietap's `CacheMap` instead. That is a left-right map: its
|
||||||
|
* `entries`/`keys`/`values` hand out live views of one of its two inner `HashMap`s and
|
||||||
|
* release the read guard before the caller iterates them, so every bulk operation here
|
||||||
|
* (`filter`, `map`, `keys()`, even a defensive `entries.toList()`) walked a `HashMap` a
|
||||||
|
* writer could be mutating. Kotlin/Native's `HashMap` detects that and throws
|
||||||
|
* `ConcurrentModificationException` — which is how `NostrClient`'s pool scans failed on
|
||||||
|
* the iOS simulator while a bunker round-trip subscribed and published from other threads.
|
||||||
|
*
|
||||||
* Iteration is weakly consistent and in bucket order. JVM/Android iterates in
|
* Iteration is weakly consistent and in bucket order. JVM/Android iterates in
|
||||||
* sorted-key order (`ConcurrentSkipListMap`) and Apple in hash order; nothing in the
|
* sorted-key order (`ConcurrentSkipListMap`); nothing in the
|
||||||
* codebase depends on a specific one. The `from`/`to` range overloads degrade to a full
|
* 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
|
* scan here, as they always have — they have no callers outside the JVM-only
|
||||||
* `LargeSoftCache`.
|
* `LargeSoftCache`.
|
||||||
+3
-1
@@ -43,7 +43,9 @@ import kotlin.concurrent.atomics.ExperimentalAtomicApi
|
|||||||
* - **One lock around a `HashMap`** writes fast but has to hand bulk operations an O(n)
|
* - **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
|
* 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`
|
* linux `PlatformLock` is a spin lock and is not reentrant, and `LocalCache`
|
||||||
* predicates reach back into the cache).
|
* predicates reach back into the cache). Apple's former `CacheMap` (left-right over
|
||||||
|
* two `HashMap`s) skipped the copy and so let scans race writers into a
|
||||||
|
* `ConcurrentModificationException`.
|
||||||
*
|
*
|
||||||
* A chained table avoids all three. Structure is only touched when a key is *added*
|
* 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
|
* (one node, prepended), an overwrite is a single volatile store into the existing
|
||||||
+39
-5
@@ -27,11 +27,12 @@ import kotlin.test.Test
|
|||||||
import kotlin.test.assertEquals
|
import kotlin.test.assertEquals
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* The linux actual is the only [LargeCache] whose thread safety is hand-rolled rather
|
* The native actual (Apple and Linux) is the only [LargeCache] whose thread safety is
|
||||||
* than delegated to a concurrent map, so it gets its own multi-threaded test. The
|
* hand-rolled rather than delegated to a concurrent map, so it gets its own
|
||||||
* copy-on-write version this replaced would fail every assertion here: its
|
* multi-threaded test. The Linux copy-on-write version it replaced would fail every
|
||||||
* read-copy-write was not a CAS loop, so concurrent writers silently dropped each
|
* assertion here: its read-copy-write was not a CAS loop, so concurrent writers silently
|
||||||
* other's entries.
|
* dropped each other's entries. The Apple `CacheMap` version it replaced failed
|
||||||
|
* [scansNeverThrowWhileKeysComeAndGo] with a `ConcurrentModificationException`.
|
||||||
*
|
*
|
||||||
* Uses `Worker` rather than coroutines on purpose — a coroutine dispatcher gives no
|
* Uses `Worker` rather than coroutines on purpose — a coroutine dispatcher gives no
|
||||||
* guarantee of genuine parallelism, and parallelism is the whole point.
|
* guarantee of genuine parallelism, and parallelism is the whole point.
|
||||||
@@ -117,4 +118,37 @@ class LargeCacheConcurrencyTest {
|
|||||||
|
|
||||||
assertEquals(1_000 + (workerCount / 2) * perWorker, cache.size())
|
assertEquals(1_000 + (workerCount / 2) * perWorker, cache.size())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun scansNeverThrowWhileKeysComeAndGo() {
|
||||||
|
val cache = LargeCache<Int, Int>()
|
||||||
|
repeat(1_000) { cache.put(it, it) }
|
||||||
|
|
||||||
|
// The shape NostrClient's pool has: one thread adds and drops subscriptions while
|
||||||
|
// another walks the whole map to rebuild filters. Every bulk read must tolerate a
|
||||||
|
// structural change landing mid-walk, removals included.
|
||||||
|
inParallel { id ->
|
||||||
|
if (id % 2 == 0) {
|
||||||
|
repeat(perWorker) { i ->
|
||||||
|
val key = 1_000 + id * perWorker + i
|
||||||
|
cache.put(key, i)
|
||||||
|
cache.remove(key - 1)
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
repeat(200) {
|
||||||
|
cache.filter { _, v -> v >= 0 }
|
||||||
|
cache.map { k, _ -> k }
|
||||||
|
cache.mapNotNull { k, v -> if (v % 2 == 0) k else null }
|
||||||
|
cache.keys().forEach { check(it >= 0) }
|
||||||
|
cache.values().forEach { check(it >= 0) }
|
||||||
|
cache.forEach { k, _ -> check(k >= 0) }
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Each writer leaves only its last key behind; its first remove takes out 999
|
||||||
|
// (the first writer's) or a key that was never there (the others').
|
||||||
|
assertEquals(999 + workerCount / 2, cache.size())
|
||||||
|
assertEquals(998, cache.get(998))
|
||||||
|
}
|
||||||
}
|
}
|
||||||
Reference in New Issue
Block a user