mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
Merge pull request #4036 from vitorpamplona/claude/linuxx64-localcache-alternative-kj1gzp
Replace copy-on-write LargeCache with striped-lock hash table
This commit is contained in:
@@ -194,6 +194,64 @@ jobs:
|
||||
name: geode Test Reports
|
||||
path: geode/build/reports
|
||||
|
||||
# Until this job existed nothing ran the linuxX64 target at all — it was compiled by
|
||||
# no CI leg. 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, and how TestResourceLoader stayed a TODO() that failed every vector-driven
|
||||
# suite on the target.
|
||||
#
|
||||
# Runs the whole :quartz suite on a Linux Native frontend, which also catches a
|
||||
# commonMain or commonTest source reaching for a JVM-only API on a target that, unlike
|
||||
# Apple, has no Foundation to fall back on.
|
||||
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 on Linux Native
|
||||
run: ./gradlew :quartz:linuxX64Test
|
||||
|
||||
- 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
|
||||
|
||||
@@ -55,7 +55,6 @@ media3 = "1.11.0"
|
||||
mockk = "1.14.11"
|
||||
kotlinx-coroutines-test = "1.11.0"
|
||||
negentropyKmp = "v1.2.0"
|
||||
netUrlencoderLibVersion = "1.6.0"
|
||||
navigationCompose = "2.9.8"
|
||||
okhttp = "5.5.0"
|
||||
osmdroid = "6.1.20"
|
||||
@@ -209,7 +208,6 @@ mockk = { group = "io.mockk", name = "mockk", version.ref = "mockk" }
|
||||
mockk-android = { group = "io.mockk", name = "mockk-android", version.ref = "mockk" }
|
||||
kotlinx-coroutines-test = { group = "org.jetbrains.kotlinx", name = "kotlinx-coroutines-test", version.ref = "kotlinx-coroutines-test"}
|
||||
negentropy-kmp = { module = "com.vitorpamplona.negentropy:kmp-negentropy", version.ref = "negentropyKmp" }
|
||||
net-thauvin-erik-urlencoder-lib = { module = "net.thauvin.erik.urlencoder:urlencoder-lib", version.ref = "netUrlencoderLibVersion" }
|
||||
okhttp = { group = "com.squareup.okhttp3", name = "okhttp", version.ref = "okhttp" }
|
||||
okhttpCoroutines = { group = "com.squareup.okhttp3", name = "okhttp-coroutines", version.ref = "okhttp" }
|
||||
osmdroid-android = { group = "org.osmdroid", name = "osmdroid-android", version.ref = "osmdroid" }
|
||||
|
||||
@@ -289,7 +289,6 @@ kotlin {
|
||||
dependsOn(nativeMain)
|
||||
dependencies {
|
||||
implementation(libs.charlietap.cachemap)
|
||||
implementation(libs.net.thauvin.erik.urlencoder.lib)
|
||||
implementation(libs.dev.whyoleg.cryptography.provider.apple.optimal)
|
||||
implementation("io.github.andreypfau:kotlinx-crypto-hmac:0.0.4")
|
||||
implementation("io.github.andreypfau:kotlinx-crypto-sha2:0.0.4")
|
||||
@@ -347,7 +346,6 @@ kotlin {
|
||||
create("linuxMain") {
|
||||
dependsOn(nativeMain)
|
||||
dependencies {
|
||||
implementation(libs.net.thauvin.erik.urlencoder.lib)
|
||||
implementation(libs.dev.whyoleg.cryptography.provider.apple.optimal)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,29 +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
|
||||
|
||||
import net.thauvin.erik.urlencoder.UrlEncoderUtil
|
||||
|
||||
actual object UrlEncoder {
|
||||
actual fun encode(value: String): String = UrlEncoderUtil.encode(value)
|
||||
|
||||
actual fun decode(value: String): String = UrlEncoderUtil.decode(value)
|
||||
}
|
||||
@@ -0,0 +1,134 @@
|
||||
/*
|
||||
* 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
|
||||
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
|
||||
/**
|
||||
* Cross-target contract for [UrlEncoder].
|
||||
*
|
||||
* This is not a style preference — [UrlEncoder.encode] builds strings that leave the
|
||||
* device. `TorrentEvent` puts it in magnet links, `Nip54InlineMetadata` in inline
|
||||
* metadata, and `Nip47DeepLink` in the `callback`/`appname`/`value` parameters of NWC
|
||||
* deep links. If Android encodes a title one way and iOS another, the two clients emit
|
||||
* different bytes for the same event, and a wallet that round-trips a deep link built
|
||||
* on one platform can fail on the other.
|
||||
*
|
||||
* The JVM/Android actual is `java.net.URLEncoder`/`URLDecoder` with UTF-8, so that is
|
||||
* the reference every other target has to match. The expectations below are its
|
||||
* `application/x-www-form-urlencoded` rules, which differ from plain RFC 3986
|
||||
* percent-encoding in exactly the three places a generic library gets "wrong": space,
|
||||
* `*` and `~`.
|
||||
*/
|
||||
class UrlEncoderTest {
|
||||
@Test
|
||||
fun keepsTheUnreservedSet() {
|
||||
// URLEncoder's dontNeedEncoding set is alphanumerics plus these four, and only
|
||||
// these four. Note `*` survives and `~` does not — the opposite of RFC 3986.
|
||||
assertEquals("abcXYZ019", UrlEncoder.encode("abcXYZ019"))
|
||||
assertEquals("-_.*", UrlEncoder.encode("-_.*"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun encodesSpaceAsPlus() {
|
||||
// Form encoding, not %20.
|
||||
assertEquals("hello+world", UrlEncoder.encode("hello world"))
|
||||
assertEquals("a+b+c", UrlEncoder.encode("a b c"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun encodesTildeAndTheOtherSubDelimiters() {
|
||||
assertEquals("%7E", UrlEncoder.encode("~"))
|
||||
assertEquals("%21", UrlEncoder.encode("!"))
|
||||
assertEquals("%27", UrlEncoder.encode("'"))
|
||||
assertEquals("%28%29", UrlEncoder.encode("()"))
|
||||
assertEquals("%24%2C%3B", UrlEncoder.encode("$,;"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun encodesUriPunctuationWithUppercaseHex() {
|
||||
assertEquals("%3A%2F%2F", UrlEncoder.encode("://"))
|
||||
assertEquals("%2B", UrlEncoder.encode("+"))
|
||||
assertEquals("%3F%26%3D", UrlEncoder.encode("?&="))
|
||||
assertEquals("%40%23%25", UrlEncoder.encode("@#%"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun encodesNonAsciiAsUtf8() {
|
||||
assertEquals("%C3%A9", UrlEncoder.encode("é"))
|
||||
assertEquals("caf%C3%A9", UrlEncoder.encode("café"))
|
||||
assertEquals("%E2%82%AC", UrlEncoder.encode("€"))
|
||||
// Outside the BMP: a surrogate pair has to encode as one 4-byte sequence.
|
||||
assertEquals("%F0%9F%98%80", UrlEncoder.encode("😀"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun decodesPlusAsSpace() {
|
||||
assertEquals("hello world", UrlEncoder.decode("hello+world"))
|
||||
assertEquals("hello world", UrlEncoder.decode("hello%20world"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun decodesPercentEscapes() {
|
||||
assertEquals("://", UrlEncoder.decode("%3A%2F%2F"))
|
||||
assertEquals("~", UrlEncoder.decode("%7E"))
|
||||
assertEquals("café", UrlEncoder.decode("caf%C3%A9"))
|
||||
assertEquals("€", UrlEncoder.decode("%E2%82%AC"))
|
||||
assertEquals("😀", UrlEncoder.decode("%F0%9F%98%80"))
|
||||
// Lowercase hex decodes the same as uppercase.
|
||||
assertEquals("é", UrlEncoder.decode("%c3%a9"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun leavesUnescapedTextAlone() {
|
||||
assertEquals("plain", UrlEncoder.decode("plain"))
|
||||
assertEquals("-_.*~", UrlEncoder.decode("-_.*~"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun roundTripsTheStringsThisIsActuallyUsedFor() {
|
||||
// A torrent title (TorrentEvent) and a tracker URL.
|
||||
val title = "Big Buck Bunny (2008) [1080p] ~ 60% done!"
|
||||
assertEquals(title, UrlEncoder.decode(UrlEncoder.encode(title)))
|
||||
|
||||
val tracker = "udp://tracker.example.org:1337/announce"
|
||||
assertEquals("udp%3A%2F%2Ftracker.example.org%3A1337%2Fannounce", UrlEncoder.encode(tracker))
|
||||
assertEquals(tracker, UrlEncoder.decode(UrlEncoder.encode(tracker)))
|
||||
|
||||
// An NWC pairing code (Nip47DeepLink.buildCallbackUri puts this in `value=`).
|
||||
val pairing =
|
||||
"nostr+walletconnect://b889ff5b1513b641e2a139f661a661364979c5beee91842f8f0ef42ab558e9d4" +
|
||||
"?relay=wss%3A%2F%2Frelay.damus.io&secret=71a8c14c1407c113601079c4302dab36460f0ccd0ad506f1f2dc73b5100571c5"
|
||||
assertEquals(pairing, UrlEncoder.decode(UrlEncoder.encode(pairing)))
|
||||
|
||||
// A callback deep link (Nip47DeepLink.parseConnectUri reads this back).
|
||||
val callback = "amethystnwc://callback"
|
||||
assertEquals("amethystnwc%3A%2F%2Fcallback", UrlEncoder.encode(callback))
|
||||
assertEquals(callback, UrlEncoder.decode(UrlEncoder.encode(callback)))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun roundTripsEveryAsciiCharacter() {
|
||||
val ascii = (0..127).map { it.toChar() }.joinToString("")
|
||||
assertEquals(ascii, UrlEncoder.decode(UrlEncoder.encode(ascii)))
|
||||
}
|
||||
}
|
||||
+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))
|
||||
}
|
||||
}
|
||||
@@ -121,41 +121,54 @@ actual class UriParser actual constructor(
|
||||
|
||||
actual fun path(): String? = parsedPath
|
||||
|
||||
actual fun queryParameterNames(): Set<String> {
|
||||
val query = parsedQuery ?: return emptySet()
|
||||
return query
|
||||
.split('&')
|
||||
.map { param ->
|
||||
val eqIndex = param.indexOf('=')
|
||||
if (eqIndex >= 0) param.substring(0, eqIndex) else param
|
||||
}.toSet()
|
||||
}
|
||||
/**
|
||||
* Parsed once and reused, mirroring the JVM actual's lazy map — the previous version
|
||||
* re-split the entire query string on every [getQueryParameter] call.
|
||||
*
|
||||
* Decoded with [UrlEncoder.decode], which matches `URLDecoder.decode(.., "UTF-8")` —
|
||||
* what the JVM actual calls. Skipping this is why NIP-47 failed on this target:
|
||||
* `relay=wss%3A%2F%2Frelay.damus.io` reached `RelayUrlNormalizer` still encoded and
|
||||
* came back "Invalid relay Url".
|
||||
*/
|
||||
private val queryParameters: Map<String, List<String>> by lazy {
|
||||
parsedQuery?.ifBlank { null }?.let { query ->
|
||||
val params = mutableMapOf<String, MutableList<String>>()
|
||||
|
||||
actual fun getQueryParameter(param: String): List<String>? {
|
||||
val query = parsedQuery ?: return null
|
||||
return query
|
||||
.split('&')
|
||||
.filter { part ->
|
||||
val eqIndex = part.indexOf('=')
|
||||
if (eqIndex >= 0) part.substring(0, eqIndex) == param else part == param
|
||||
}.map { part ->
|
||||
val eqIndex = part.indexOf('=')
|
||||
if (eqIndex >= 0) part.substring(eqIndex + 1) else ""
|
||||
query.split('&').forEach { paramValue ->
|
||||
val parts = paramValue.split("=", limit = 2)
|
||||
val currentValue =
|
||||
params.getOrPut(parts[0]) {
|
||||
mutableListOf()
|
||||
}
|
||||
|
||||
if (parts.size == 2) {
|
||||
currentValue.add(UrlEncoder.decode(parts[1]))
|
||||
} else {
|
||||
currentValue.add("")
|
||||
}
|
||||
}
|
||||
|
||||
params
|
||||
} ?: emptyMap()
|
||||
}
|
||||
|
||||
val fragments: Map<String, String> by lazy {
|
||||
private val parsedFragments: Map<String, String> by lazy {
|
||||
parsedFragment?.ifBlank { null }?.let { keyValuePair ->
|
||||
keyValuePair.split('&').associate { paramValue ->
|
||||
val parts = paramValue.split("=", limit = 2)
|
||||
if (parts.size == 2) {
|
||||
parts[0] to parts[1]
|
||||
parts[0] to UrlEncoder.decode(parts[1])
|
||||
} else {
|
||||
parts[0] to ""
|
||||
parts[0] to "" // Handle parameters without a value
|
||||
}
|
||||
}
|
||||
} ?: emptyMap()
|
||||
}
|
||||
|
||||
actual fun fragments(): Map<String, String> = fragments
|
||||
actual fun queryParameterNames(): Set<String> = queryParameters.keys
|
||||
|
||||
/** Null — not an empty list — when the parameter is absent, as on the JVM. */
|
||||
actual fun getQueryParameter(param: String): List<String>? = queryParameters[param]
|
||||
|
||||
actual fun fragments(): Map<String, String> = parsedFragments
|
||||
}
|
||||
|
||||
@@ -1,29 +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
|
||||
|
||||
import net.thauvin.erik.urlencoder.UrlEncoderUtil
|
||||
|
||||
actual object UrlEncoder {
|
||||
actual fun encode(value: String): String = UrlEncoderUtil.encode(value)
|
||||
|
||||
actual fun decode(value: String): String = UrlEncoderUtil.decode(value)
|
||||
}
|
||||
Vendored
+14
-14
@@ -20,30 +20,30 @@
|
||||
*/
|
||||
package com.vitorpamplona.quartz.utils.cache
|
||||
|
||||
import kotlin.concurrent.AtomicReference
|
||||
|
||||
// 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], over the same [StripedHashMap] as
|
||||
* `LargeCache.linux.kt` — read its docs for why.
|
||||
*
|
||||
* This one was the worst-placed of the copy-on-write caches: its only caller,
|
||||
* `CachingEventDecoder`, writes once per event arriving from a relay, so every decode
|
||||
* rebuilt the whole map under a CAS retry loop. Now a write touches one bucket and a
|
||||
* read takes no lock at all.
|
||||
*/
|
||||
actual class ConcurrentHashCache<K : Any, V : Any> {
|
||||
private val mapRef = AtomicReference(HashMap<K, V>())
|
||||
private val cache = StripedHashMap<K, V>()
|
||||
|
||||
actual fun get(key: K): V? = mapRef.value[key]
|
||||
actual fun get(key: K): V? = cache.get(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
|
||||
}
|
||||
cache.put(key, value)
|
||||
}
|
||||
|
||||
actual fun size(): Int = mapRef.value.size
|
||||
actual fun size(): Int = cache.size()
|
||||
|
||||
actual fun clear() {
|
||||
mapRef.value = HashMap()
|
||||
cache.clear()
|
||||
}
|
||||
}
|
||||
|
||||
+159
-108
@@ -20,175 +20,226 @@
|
||||
*/
|
||||
package com.vitorpamplona.quartz.utils.cache
|
||||
|
||||
import kotlin.concurrent.AtomicReference
|
||||
|
||||
/**
|
||||
* Linux/Native actual for [LargeCache] — the store behind Amethyst's `LocalCache`.
|
||||
*
|
||||
* All of the concurrency and the performance rationale lives in [StripedHashMap]; this
|
||||
* class is only the [ICacheOperations] surface over it. The short version: `LocalCache`
|
||||
* fills ~100,000 entries in a few seconds while feeds scan the whole cache, so the
|
||||
* backing store has to take writes at O(1) with next to no garbage *and* let a scan run
|
||||
* in place without copying. A chained hash table with lock-free reads and striped-lock
|
||||
* writes — `ConcurrentHashMap`'s shape, which Kotlin/Native does not ship — is the
|
||||
* structure that does both; copy-on-write and a persistent HAMT each fail one half.
|
||||
*
|
||||
* Every bulk operation below walks the table through [StripedHashMap.forEachEntry],
|
||||
* which takes no lock and allocates nothing beyond the result being built. Two things
|
||||
* follow, both of which the earlier implementations had to work around:
|
||||
*
|
||||
* - The caller's lambda never runs inside a critical section, so a `LocalCache`
|
||||
* predicate that reaches back into the cache cannot deadlock.
|
||||
* - There is no snapshot and no defensive `entries.toList()`, so no
|
||||
* `ConcurrentModificationException` window and no per-scan copy.
|
||||
*
|
||||
* Iteration is weakly consistent and in bucket order. JVM/Android iterates in
|
||||
* sorted-key order (`ConcurrentSkipListMap`) and Apple in hash order; nothing in the
|
||||
* codebase depends on a specific one. The `from`/`to` range overloads degrade to a full
|
||||
* scan here, as they always have — they have no callers outside the JVM-only
|
||||
* `LargeSoftCache`.
|
||||
*/
|
||||
actual class LargeCache<K, V> : ICacheOperations<K, V> {
|
||||
private val mapRef = AtomicReference(LinkedHashMap<K, V>())
|
||||
private val cache = StripedHashMap<K, V>()
|
||||
|
||||
private inline fun <R> withMap(block: (LinkedHashMap<K, V>) -> R): R = block(mapRef.value)
|
||||
|
||||
private inline fun mutate(block: (LinkedHashMap<K, V>) -> Unit) {
|
||||
val copy = LinkedHashMap(mapRef.value)
|
||||
block(copy)
|
||||
mapRef.value = copy
|
||||
actual fun keys(): Set<K> {
|
||||
val results = LinkedHashSet<K>(cache.size())
|
||||
cache.forEachEntry { key, _ -> results.add(key) }
|
||||
return results
|
||||
}
|
||||
|
||||
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) }
|
||||
}
|
||||
return value
|
||||
actual fun values(): Iterable<V> {
|
||||
val results = ArrayList<V>(cache.size())
|
||||
cache.forEachEntry { _, value -> results.add(value) }
|
||||
return results
|
||||
}
|
||||
|
||||
actual fun isEmpty(): Boolean = withMap { it.isEmpty() }
|
||||
actual fun get(key: K): V? = cache.get(key)
|
||||
|
||||
actual fun remove(key: K): V? = cache.remove(key)
|
||||
|
||||
actual fun isEmpty(): Boolean = cache.isEmpty()
|
||||
|
||||
actual fun clear() {
|
||||
mapRef.value = LinkedHashMap()
|
||||
cache.clear()
|
||||
}
|
||||
|
||||
actual fun containsKey(key: K): Boolean = withMap { it.containsKey(key) }
|
||||
actual fun containsKey(key: K): Boolean = cache.containsKey(key)
|
||||
|
||||
actual fun put(
|
||||
key: K,
|
||||
value: V,
|
||||
) {
|
||||
mutate { it[key] = value }
|
||||
cache.put(key, value)
|
||||
}
|
||||
|
||||
// The next two are the JVM actual's bodies verbatim, over the same putIfAbsent
|
||||
// contract: [builder] runs outside the write path, and the loser of a race keeps
|
||||
// the winner's value.
|
||||
|
||||
actual fun getOrCreate(
|
||||
key: K,
|
||||
builder: (key: K) -> V,
|
||||
): V {
|
||||
val existing = get(key)
|
||||
if (existing != null) return existing
|
||||
val newObject = builder(key)
|
||||
mutate { it[key] = newObject }
|
||||
return get(key) ?: newObject
|
||||
val value = cache.get(key)
|
||||
|
||||
return if (value != null) {
|
||||
value
|
||||
} else {
|
||||
val newObject = builder(key)
|
||||
cache.putIfAbsent(key, newObject) ?: newObject
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* True only when *this* call inserted. An early implementation returned
|
||||
* `get(key) != null`, which also reported true when another thread had just created
|
||||
* the entry, double-firing whatever the caller does with a fresh key.
|
||||
*/
|
||||
actual fun createIfAbsent(
|
||||
key: K,
|
||||
builder: (key: K) -> V,
|
||||
): Boolean {
|
||||
val existing = get(key)
|
||||
if (existing != null) return false
|
||||
val newObject = builder(key)
|
||||
mutate { it[key] = newObject }
|
||||
return get(key) != null
|
||||
val value = cache.get(key)
|
||||
return if (value != null) {
|
||||
false
|
||||
} else {
|
||||
val newObject = builder(key)
|
||||
cache.putIfAbsent(key, newObject) == null
|
||||
}
|
||||
}
|
||||
|
||||
actual override fun size(): Int = withMap { it.size }
|
||||
actual override fun size(): Int = cache.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) }
|
||||
cache.forEachEntry { key, value -> consumer.accept(key, value) }
|
||||
}
|
||||
|
||||
actual override fun filter(consumer: CacheCollectors.BiFilter<K, V>): List<V> = withMap { map -> map.filter { consumer.filter(it.key, it.value) }.values.toList() }
|
||||
actual override fun filter(consumer: CacheCollectors.BiFilter<K, V>): List<V> {
|
||||
val results = ArrayList<V>()
|
||||
cache.forEachEntry { key, value -> if (consumer.filter(key, value)) results.add(value) }
|
||||
return results
|
||||
}
|
||||
|
||||
actual override fun filterIntoSet(consumer: CacheCollectors.BiFilter<K, V>): Set<V> = withMap { map -> map.filter { consumer.filter(it.key, it.value) }.values.toSet() }
|
||||
actual override fun filterIntoSet(consumer: CacheCollectors.BiFilter<K, V>): Set<V> {
|
||||
val results = LinkedHashSet<V>()
|
||||
cache.forEachEntry { key, value -> if (consumer.filter(key, value)) results.add(value) }
|
||||
return results
|
||||
}
|
||||
|
||||
actual override fun <R> map(consumer: CacheCollectors.BiNotNullMapper<K, V, R>): List<R> = withMap { map -> map.map { consumer.map(it.key, it.value) } }
|
||||
actual override fun <R> map(consumer: CacheCollectors.BiNotNullMapper<K, V, R>): List<R> {
|
||||
val results = ArrayList<R>(cache.size())
|
||||
cache.forEachEntry { key, value -> results.add(consumer.map(key, value)) }
|
||||
return results
|
||||
}
|
||||
|
||||
actual override fun <R> mapNotNull(consumer: CacheCollectors.BiMapper<K, V, R?>): List<R> = withMap { map -> map.mapNotNull { consumer.map(it.key, it.value) } }
|
||||
actual override fun <R> mapNotNull(consumer: CacheCollectors.BiMapper<K, V, R?>): List<R> {
|
||||
val results = ArrayList<R>()
|
||||
cache.forEachEntry { key, value -> consumer.map(key, value)?.let { results.add(it) } }
|
||||
return results
|
||||
}
|
||||
|
||||
actual override fun <R> mapNotNullIntoSet(consumer: CacheCollectors.BiMapper<K, V, R?>): Set<R> = mapNotNull(consumer).toSet()
|
||||
actual override fun <R> mapNotNullIntoSet(consumer: CacheCollectors.BiMapper<K, V, R?>): Set<R> {
|
||||
val results = LinkedHashSet<R>()
|
||||
cache.forEachEntry { key, value -> consumer.map(key, value)?.let { results.add(it) } }
|
||||
return results
|
||||
}
|
||||
|
||||
actual override fun <R> mapFlatten(consumer: CacheCollectors.BiMapper<K, V, Collection<R>?>): List<R> = withMap { map -> map.flatMap { entry -> consumer.map(entry.key, entry.value) ?: emptyList() } }
|
||||
actual override fun <R> mapFlatten(consumer: CacheCollectors.BiMapper<K, V, Collection<R>?>): List<R> {
|
||||
val results = ArrayList<R>()
|
||||
cache.forEachEntry { key, value -> consumer.map(key, value)?.let { results.addAll(it) } }
|
||||
return results
|
||||
}
|
||||
|
||||
actual override fun <R> mapFlattenIntoSet(consumer: CacheCollectors.BiMapper<K, V, Collection<R>?>): Set<R> = mapFlatten(consumer).toSet()
|
||||
actual override fun <R> mapFlattenIntoSet(consumer: CacheCollectors.BiMapper<K, V, Collection<R>?>): Set<R> {
|
||||
val results = LinkedHashSet<R>()
|
||||
cache.forEachEntry { key, value -> consumer.map(key, value)?.let { results.addAll(it) } }
|
||||
return results
|
||||
}
|
||||
|
||||
actual override fun maxOrNullOf(
|
||||
filter: CacheCollectors.BiFilter<K, V>,
|
||||
comparator: Comparator<V>,
|
||||
): V? =
|
||||
withMap { map ->
|
||||
var maxV: V? = null
|
||||
map.forEach {
|
||||
if (filter.filter(it.key, it.value)) {
|
||||
if (maxV == null || comparator.compare(it.value, maxV) > 0) {
|
||||
maxV = it.value
|
||||
}
|
||||
): V? {
|
||||
var maxV: V? = null
|
||||
cache.forEachEntry { key, value ->
|
||||
if (filter.filter(key, value)) {
|
||||
if (maxV == null || comparator.compare(value, maxV) > 0) {
|
||||
maxV = value
|
||||
}
|
||||
}
|
||||
maxV
|
||||
}
|
||||
return maxV
|
||||
}
|
||||
|
||||
actual override fun sumOf(consumer: CacheCollectors.BiSumOf<K, V>): Int =
|
||||
withMap { map ->
|
||||
var sum = 0
|
||||
map.forEach { sum += consumer.map(it.key, it.value) }
|
||||
sum
|
||||
}
|
||||
actual override fun sumOf(consumer: CacheCollectors.BiSumOf<K, V>): Int {
|
||||
var sum = 0
|
||||
cache.forEachEntry { key, value -> sum += consumer.map(key, value) }
|
||||
return sum
|
||||
}
|
||||
|
||||
actual override fun sumOfLong(consumer: CacheCollectors.BiSumOfLong<K, V>): Long =
|
||||
withMap { map ->
|
||||
var sum = 0L
|
||||
map.forEach { sum += consumer.map(it.key, it.value) }
|
||||
sum
|
||||
}
|
||||
actual override fun sumOfLong(consumer: CacheCollectors.BiSumOfLong<K, V>): Long {
|
||||
var sum = 0L
|
||||
cache.forEachEntry { key, value -> sum += consumer.map(key, value) }
|
||||
return sum
|
||||
}
|
||||
|
||||
actual override fun <R> groupBy(consumer: CacheCollectors.BiNotNullMapper<K, V, R>): Map<R, List<V>> =
|
||||
withMap { map ->
|
||||
val results = HashMap<R, ArrayList<V>>()
|
||||
map.forEach {
|
||||
val group = consumer.map(it.key, it.value)
|
||||
results.getOrPut(group) { ArrayList() }.add(it.value)
|
||||
}
|
||||
results
|
||||
actual override fun <R> groupBy(consumer: CacheCollectors.BiNotNullMapper<K, V, R>): Map<R, List<V>> {
|
||||
val results = HashMap<R, ArrayList<V>>()
|
||||
cache.forEachEntry { key, value ->
|
||||
results.getOrPut(consumer.map(key, value)) { ArrayList() }.add(value)
|
||||
}
|
||||
return results
|
||||
}
|
||||
|
||||
actual override fun <R> countByGroup(consumer: CacheCollectors.BiNotNullMapper<K, V, R>): Map<R, Int> =
|
||||
withMap { map ->
|
||||
val results = HashMap<R, Int>()
|
||||
map.forEach {
|
||||
val group = consumer.map(it.key, it.value)
|
||||
results[group] = (results[group] ?: 0) + 1
|
||||
}
|
||||
results
|
||||
actual override fun <R> countByGroup(consumer: CacheCollectors.BiNotNullMapper<K, V, R>): Map<R, Int> {
|
||||
val results = HashMap<R, Int>()
|
||||
cache.forEachEntry { key, value ->
|
||||
val group = consumer.map(key, value)
|
||||
results[group] = (results[group] ?: 0) + 1
|
||||
}
|
||||
return results
|
||||
}
|
||||
|
||||
actual override fun <R> sumByGroup(
|
||||
groupMap: CacheCollectors.BiNotNullMapper<K, V, R>,
|
||||
sumOf: CacheCollectors.BiNotNullMapper<K, V, Long>,
|
||||
): Map<R, Long> =
|
||||
withMap { map ->
|
||||
val results = HashMap<R, Long>()
|
||||
map.forEach {
|
||||
val group = groupMap.map(it.key, it.value)
|
||||
results[group] = (results[group] ?: 0L) + sumOf.map(it.key, it.value)
|
||||
}
|
||||
results
|
||||
): Map<R, Long> {
|
||||
val results = HashMap<R, Long>()
|
||||
cache.forEachEntry { key, value ->
|
||||
val group = groupMap.map(key, value)
|
||||
results[group] = (results[group] ?: 0L) + sumOf.map(key, value)
|
||||
}
|
||||
return results
|
||||
}
|
||||
|
||||
actual override fun count(consumer: CacheCollectors.BiFilter<K, V>): Int = withMap { map -> map.count { consumer.filter(it.key, it.value) } }
|
||||
actual override fun count(consumer: CacheCollectors.BiFilter<K, V>): Int {
|
||||
var count = 0
|
||||
cache.forEachEntry { key, value -> if (consumer.filter(key, value)) count++ }
|
||||
return count
|
||||
}
|
||||
|
||||
actual override fun <T, U> associate(transform: (K, V) -> Pair<T, U>): Map<T, U> =
|
||||
withMap { map ->
|
||||
val results = LinkedHashMap<T, U>(map.size)
|
||||
map.forEach {
|
||||
val pair = transform(it.key, it.value)
|
||||
results[pair.first] = pair.second
|
||||
}
|
||||
results
|
||||
actual override fun <T, U> associate(transform: (K, V) -> Pair<T, U>): Map<T, U> {
|
||||
val results = LinkedHashMap<T, U>(cache.size())
|
||||
cache.forEachEntry { key, value ->
|
||||
val pair = transform(key, value)
|
||||
results[pair.first] = pair.second
|
||||
}
|
||||
return results
|
||||
}
|
||||
|
||||
actual override fun <U> associateWith(transform: (K, V) -> U?): Map<K, U?> =
|
||||
withMap { map ->
|
||||
val results = LinkedHashMap<K, U?>(map.size)
|
||||
map.forEach {
|
||||
results[it.key] = transform(it.key, it.value)
|
||||
}
|
||||
results
|
||||
}
|
||||
actual override fun <U> associateWith(transform: (K, V) -> U?): Map<K, U?> {
|
||||
val results = LinkedHashMap<K, U?>(cache.size())
|
||||
cache.forEachEntry { key, value -> results[key] = transform(key, value) }
|
||||
return results
|
||||
}
|
||||
|
||||
actual override fun filter(
|
||||
from: K,
|
||||
|
||||
+342
@@ -0,0 +1,342 @@
|
||||
/*
|
||||
* Copyright (c) 2025 Vitor Pamplona
|
||||
*
|
||||
* Permission is hereby granted, free of charge, to any person obtaining a copy of
|
||||
* this software and associated documentation files (the "Software"), to deal in
|
||||
* the Software without restriction, including without limitation the rights to use,
|
||||
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
|
||||
* Software, and to permit persons to whom the Software is furnished to do so,
|
||||
* subject to the following conditions:
|
||||
*
|
||||
* The above copyright notice and this permission notice shall be included in all
|
||||
* copies or substantial portions of the Software.
|
||||
*
|
||||
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
|
||||
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
|
||||
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
|
||||
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
|
||||
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
|
||||
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
||||
*/
|
||||
package com.vitorpamplona.quartz.utils.cache
|
||||
|
||||
import com.vitorpamplona.quartz.utils.concurrent.PlatformLock
|
||||
import com.vitorpamplona.quartz.utils.concurrent.withLock
|
||||
import kotlin.concurrent.Volatile
|
||||
import kotlin.concurrent.atomics.AtomicArray
|
||||
import kotlin.concurrent.atomics.AtomicInt
|
||||
import kotlin.concurrent.atomics.ExperimentalAtomicApi
|
||||
|
||||
/**
|
||||
* A chained hash table with **lock-free reads** and **striped-lock writes** — the shape
|
||||
* of `java.util.concurrent.ConcurrentHashMap`, which Kotlin/Native has no equivalent of.
|
||||
*
|
||||
* Exists because `LocalCache` fills on the order of 100,000 entries in a few seconds,
|
||||
* and every structure available on this target fails that workload in some way:
|
||||
*
|
||||
* - **Copy-on-write over a `HashMap`** (what shipped first) rebuilds the whole map per
|
||||
* write: O(n) each, O(n^2) to fill.
|
||||
* - **A persistent HAMT + CAS** is O(log32 n) per write, but allocates a fresh path of
|
||||
* ~4-5 nodes for *every* write — including overwrites, which change no structure at
|
||||
* all — and throws the old path away. Measured over a 100k fill plus scans that is 24
|
||||
* GC cycles against this table's 1.
|
||||
* - **One lock around a `HashMap`** writes fast but has to hand bulk operations an O(n)
|
||||
* copy, because a caller's lambda must not run inside the critical section (the
|
||||
* linux `PlatformLock` is a spin lock and is not reentrant, and `LocalCache`
|
||||
* predicates reach back into the cache).
|
||||
*
|
||||
* A chained table avoids all three. Structure is only touched when a key is *added*
|
||||
* (one node, prepended), an overwrite is a single volatile store into the existing
|
||||
* node, and scans walk the buckets in place with no copy and no lock — so caller
|
||||
* lambdas run outside any critical section and cannot deadlock.
|
||||
*
|
||||
* Measured on linuxX64 (`-opt`), 100,000 String keys of event-id length, ms per phase
|
||||
* and GC cycles over the whole run:
|
||||
*
|
||||
* ```
|
||||
* fill overwrite reads 20 scans mixed GCs heap
|
||||
* copy-on-write* n/a n/a n/a n/a n/a n/a n/a
|
||||
* HAMT + CAS 70 78 6 117 676 24 67MB
|
||||
* lock + HashMap 13 6 3 71 1197 36 51MB
|
||||
* this 16 4 1 16 86 1 41MB
|
||||
* ```
|
||||
*
|
||||
* (*copy-on-write is off the scale: 20k entries alone took 18s to fill.) "mixed" is a
|
||||
* full fill with a whole-table scan every 1000 writes, which is the shape `LocalCache`
|
||||
* actually has — arriving events interleaved with feeds filtering the whole cache.
|
||||
* Figures are one representative run of several; they were stable to within ~10%, except
|
||||
* the HAMT's scan column, which wandered between 115ms and 190ms. This row was
|
||||
* re-measured on the shipped code, after the stripe-selection fix and after
|
||||
* [INITIAL_CAPACITY] dropped to [STRIPES] — neither moved it out of the noise.
|
||||
*
|
||||
* An *empty* instance costs ~970 bytes, nearly all of it the 16 [PlatformLock]s (two
|
||||
* objects each). That is well above the ~50 bytes an empty map used to cost, and it is
|
||||
* charged to every one of the hundreds of small caches a client holds. Folding the stripe
|
||||
* locks into a single `AtomicIntArray` would take it to ~250 bytes and is the obvious next
|
||||
* step if it ever shows up in a heap profile; it is left alone here because hand-rolling
|
||||
* the spin is exactly the kind of change that wants its own review.
|
||||
*
|
||||
* ## Concurrency contract
|
||||
*
|
||||
* - **Readers never block and never allocate.** [get], [containsKey], [size] and
|
||||
* [forEachEntry] take no lock. A reader loads [table] once and walks immutable
|
||||
* `next` links, so it always sees a well-formed chain.
|
||||
* - **Writers block only against writers hashing to the same stripe**, and only for a
|
||||
* bucket walk of a few nodes. This is where it differs from the JVM actual's fully
|
||||
* non-blocking `ConcurrentSkipListMap`; it matches `ConcurrentHashMap`, which also
|
||||
* locks a bin to write it.
|
||||
* - Iteration is **weakly consistent**, like both of those: it reflects the table as of
|
||||
* its first load and may or may not observe writes that land while it runs. It never
|
||||
* throws, never sees a torn chain, and never needs a defensive copy.
|
||||
* - A resize takes every stripe lock, so no write can be in flight while it runs.
|
||||
* Nodes are rebuilt rather than relinked, which is what lets a reader that captured
|
||||
* the pre-resize table keep walking it safely.
|
||||
*/
|
||||
@OptIn(ExperimentalAtomicApi::class)
|
||||
internal class StripedHashMap<K, V> {
|
||||
/**
|
||||
* [value] is mutable so that overwriting an existing key allocates nothing; [next]
|
||||
* is not, so that a reader walking a chain can never see it change under them.
|
||||
* Structural edits publish a new head instead.
|
||||
*/
|
||||
internal class Node<K, V>(
|
||||
val hash: Int,
|
||||
val key: K,
|
||||
@Volatile var value: V,
|
||||
val next: Node<K, V>?,
|
||||
)
|
||||
|
||||
private val locks = Array(STRIPES) { PlatformLock() }
|
||||
private val entryCount = AtomicInt(0)
|
||||
|
||||
@Volatile internal var table = AtomicArray<Node<K, V>?>(INITIAL_CAPACITY) { null }
|
||||
|
||||
@Volatile private var threshold = INITIAL_CAPACITY / 4 * 3
|
||||
|
||||
/** Spreads the high bits down, so that both the bucket and the stripe see entropy. */
|
||||
private fun hashOf(key: K): Int {
|
||||
val h = key?.hashCode() ?: 0
|
||||
return h xor (h ushr 16)
|
||||
}
|
||||
|
||||
/**
|
||||
* **The stripe must be a function of the bucket**, or the locks do not partition the
|
||||
* table and the whole design is unsound: two keys could share a bucket while holding
|
||||
* different locks, so two writers would read the same chain head and both publish
|
||||
* over it, silently dropping one insert.
|
||||
*
|
||||
* It is a function of the bucket here because [STRIPES] and every table capacity are
|
||||
* powers of two with `STRIPES <= capacity`, so `hash and (STRIPES - 1)` is exactly the
|
||||
* low bits of `hash and (capacity - 1)`. Same bucket therefore implies same stripe, at
|
||||
* every size. Using the hash and not the capacity also keeps a key on one stripe
|
||||
* across a resize, which is what lets [growTable] exclude writers by taking all of
|
||||
* them.
|
||||
*
|
||||
* An earlier version took bits 16-19 instead. That is still resize-stable and still
|
||||
* spreads well, which is why it looked right — but it is not derived from the bucket,
|
||||
* so it broke the invariant above.
|
||||
*/
|
||||
private fun lockFor(hash: Int) = locks[hash and (STRIPES - 1)]
|
||||
|
||||
fun size(): Int = entryCount.load()
|
||||
|
||||
fun isEmpty(): Boolean = entryCount.load() == 0
|
||||
|
||||
fun get(key: K): V? {
|
||||
val hash = hashOf(key)
|
||||
val current = table
|
||||
var node = current.loadAt(hash and (current.size - 1))
|
||||
while (node != null) {
|
||||
if (node.hash == hash && node.key == key) return node.value
|
||||
node = node.next
|
||||
}
|
||||
return null
|
||||
}
|
||||
|
||||
fun containsKey(key: K): Boolean {
|
||||
val hash = hashOf(key)
|
||||
val current = table
|
||||
var node = current.loadAt(hash and (current.size - 1))
|
||||
while (node != null) {
|
||||
if (node.hash == hash && node.key == key) return true
|
||||
node = node.next
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
fun put(
|
||||
key: K,
|
||||
value: V,
|
||||
) {
|
||||
val hash = hashOf(key)
|
||||
var grew = false
|
||||
lockFor(hash).withLock {
|
||||
val current = table
|
||||
val index = hash and (current.size - 1)
|
||||
val head = current.loadAt(index)
|
||||
var node = head
|
||||
while (node != null) {
|
||||
if (node.hash == hash && node.key == key) {
|
||||
// Present already: no structural change, no allocation.
|
||||
node.value = value
|
||||
return@withLock
|
||||
}
|
||||
node = node.next
|
||||
}
|
||||
current.storeAt(index, Node(hash, key, value, head))
|
||||
grew = entryCount.fetchAndAdd(1) + 1 > threshold
|
||||
}
|
||||
if (grew) growTable()
|
||||
}
|
||||
|
||||
/**
|
||||
* Inserts [value] only if [key] is absent, and returns the value already stored —
|
||||
* or null when this call performed the insert. Exactly `ConcurrentMap.putIfAbsent`,
|
||||
* which the JVM actual builds `getOrCreate` and `createIfAbsent` out of, including
|
||||
* its inability to represent a stored null (`ConcurrentSkipListMap` rejects those).
|
||||
*/
|
||||
fun putIfAbsent(
|
||||
key: K,
|
||||
value: V,
|
||||
): V? {
|
||||
val hash = hashOf(key)
|
||||
var grew = false
|
||||
var existing: V? = null
|
||||
lockFor(hash).withLock {
|
||||
val current = table
|
||||
val index = hash and (current.size - 1)
|
||||
val head = current.loadAt(index)
|
||||
var node = head
|
||||
while (node != null) {
|
||||
if (node.hash == hash && node.key == key) {
|
||||
existing = node.value
|
||||
return@withLock
|
||||
}
|
||||
node = node.next
|
||||
}
|
||||
current.storeAt(index, Node(hash, key, value, head))
|
||||
grew = entryCount.fetchAndAdd(1) + 1 > threshold
|
||||
}
|
||||
if (grew) growTable()
|
||||
return existing
|
||||
}
|
||||
|
||||
fun remove(key: K): V? {
|
||||
val hash = hashOf(key)
|
||||
var removed: V? = null
|
||||
lockFor(hash).withLock {
|
||||
val current = table
|
||||
val index = hash and (current.size - 1)
|
||||
val head = current.loadAt(index)
|
||||
|
||||
var target = head
|
||||
while (target != null && !(target.hash == hash && target.key == key)) target = target.next
|
||||
if (target == null) return@withLock
|
||||
|
||||
// `next` is immutable, so the nodes ahead of the removed one are cloned onto
|
||||
// its tail. A reader still walking the old head sees the entry one last time
|
||||
// rather than a broken chain.
|
||||
var rebuilt = target.next
|
||||
var ahead = head
|
||||
while (ahead !== target) {
|
||||
val node = ahead!!
|
||||
rebuilt = Node(node.hash, node.key, node.value, rebuilt)
|
||||
ahead = node.next
|
||||
}
|
||||
|
||||
current.storeAt(index, rebuilt)
|
||||
entryCount.fetchAndAdd(-1)
|
||||
removed = target.value
|
||||
}
|
||||
return removed
|
||||
}
|
||||
|
||||
fun clear() {
|
||||
lockAll()
|
||||
try {
|
||||
table = AtomicArray(INITIAL_CAPACITY) { null }
|
||||
threshold = INITIAL_CAPACITY / 4 * 3
|
||||
entryCount.store(0)
|
||||
} finally {
|
||||
unlockAll()
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Walks every entry without locking. Inline so the caller's body runs with no
|
||||
* `Function2` dispatch and no captured-variable box per entry, which is what keeps
|
||||
* a full-cache scan allocation-free.
|
||||
*/
|
||||
inline fun forEachEntry(action: (K, V) -> Unit) {
|
||||
val current = table
|
||||
for (index in 0 until current.size) {
|
||||
var node = current.loadAt(index)
|
||||
while (node != null) {
|
||||
action(node.key, node.value)
|
||||
node = node.next
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private fun growTable() {
|
||||
lockAll()
|
||||
try {
|
||||
val old = table
|
||||
// Another writer may have grown it while this one waited for the locks.
|
||||
if (entryCount.load() <= threshold) return
|
||||
if (old.size >= MAX_CAPACITY) {
|
||||
threshold = Int.MAX_VALUE
|
||||
return
|
||||
}
|
||||
|
||||
val capacity = old.size shl 1
|
||||
val next = AtomicArray<Node<K, V>?>(capacity) { null }
|
||||
for (index in 0 until old.size) {
|
||||
var node = old.loadAt(index)
|
||||
while (node != null) {
|
||||
val target = node.hash and (capacity - 1)
|
||||
next.storeAt(target, Node(node.hash, node.key, node.value, next.loadAt(target)))
|
||||
node = node.next
|
||||
}
|
||||
}
|
||||
|
||||
table = next
|
||||
threshold = capacity / 4 * 3
|
||||
} finally {
|
||||
unlockAll()
|
||||
}
|
||||
}
|
||||
|
||||
/** Always in index order, and only ever from a thread holding no stripe lock. */
|
||||
private fun lockAll() {
|
||||
for (lock in locks) lock.lock()
|
||||
}
|
||||
|
||||
private fun unlockAll() {
|
||||
for (index in locks.indices.reversed()) locks[index].unlock()
|
||||
}
|
||||
|
||||
companion object {
|
||||
/**
|
||||
* Writes are O(1), so a stripe is held for a few nanoseconds and 16 ways is
|
||||
* plenty — `ConcurrentHashMap` shipped with the same default for years.
|
||||
*/
|
||||
private const val STRIPES = 16
|
||||
|
||||
/**
|
||||
* Tied to [STRIPES] rather than chosen independently: the stripe is only a
|
||||
* function of the bucket while `STRIPES <= capacity` (see [lockFor]), and defining
|
||||
* it this way makes that impossible to break by raising [STRIPES] alone.
|
||||
*
|
||||
* Kept small on purpose. `LargeCache` is not only the one big `LocalCache`
|
||||
* instance: `EphemeralRoom`, `RelaySession`, `PoolRequests` and friends each build
|
||||
* one per room, per connection and per subscription set, so a client holds
|
||||
* hundreds of them and most stay nearly empty. Starting at 1024 slots charged every
|
||||
* one of those ~8 KB it would never use. Growth is geometric, so a table that does
|
||||
* fill to 100k pays the same ~2n node rebuilds in total either way.
|
||||
*/
|
||||
private const val INITIAL_CAPACITY = STRIPES
|
||||
|
||||
private const val MAX_CAPACITY = 1 shl 30
|
||||
}
|
||||
}
|
||||
@@ -20,12 +20,82 @@
|
||||
*/
|
||||
package com.vitorpamplona.quartz
|
||||
|
||||
actual class TestResourceLoader actual constructor() {
|
||||
actual fun loadDecompressString(file: String): String {
|
||||
TODO("Not yet implemented")
|
||||
}
|
||||
import com.vitorpamplona.quartz.utils.GZip
|
||||
import kotlinx.cinterop.ExperimentalForeignApi
|
||||
import kotlinx.cinterop.addressOf
|
||||
import kotlinx.cinterop.convert
|
||||
import kotlinx.cinterop.toKString
|
||||
import kotlinx.cinterop.usePinned
|
||||
import platform.posix.SEEK_END
|
||||
import platform.posix.SEEK_SET
|
||||
import platform.posix.fclose
|
||||
import platform.posix.fopen
|
||||
import platform.posix.fread
|
||||
import platform.posix.fseek
|
||||
import platform.posix.ftell
|
||||
import platform.posix.getenv
|
||||
|
||||
actual fun loadString(file: String): String {
|
||||
TODO("Not yet implemented")
|
||||
/**
|
||||
* Linux/Native actual for [TestResourceLoader].
|
||||
*
|
||||
* Until this existed it was `TODO()`, which failed 69 tests on this target — every
|
||||
* suite driven by a vector file: the whole MLS interop set, NIP-44, the NIP-01 hint
|
||||
* indexer, the SQLite store's large-DB tests and the Bolt12 payer proofs. None of them
|
||||
* were failing because of missing production code; they could not read their input.
|
||||
*
|
||||
* Resolves paths against `TEST_RESOURCES_ROOT`, the same environment variable the Apple
|
||||
* actual uses, exported onto every `KotlinNativeTest` task by `quartz/build.gradle.kts`.
|
||||
* Reads through `platform.posix` rather than Foundation, which linuxX64 does not have.
|
||||
*
|
||||
* The read is a single `stat`-sized allocation filled by `fread`, so a vector file
|
||||
* costs exactly one `ByteArray` — less than the JVM actual's `bufferedReader().readText()`,
|
||||
* which grows a `StringBuilder` as it goes.
|
||||
*/
|
||||
@OptIn(ExperimentalForeignApi::class)
|
||||
actual class TestResourceLoader actual constructor() {
|
||||
actual fun loadDecompressString(file: String): String = GZip.decompress(readBytes(file))
|
||||
|
||||
actual fun loadString(file: String): String = readBytes(file).decodeToString()
|
||||
|
||||
private fun readBytes(file: String): ByteArray {
|
||||
val root =
|
||||
getenv("TEST_RESOURCES_ROOT")?.toKString()
|
||||
?: throw IllegalStateException(
|
||||
"TEST_RESOURCES_ROOT is not set. quartz/build.gradle.kts exports it onto every " +
|
||||
"KotlinNativeTest task; running the test binary directly has to set it too.",
|
||||
)
|
||||
|
||||
val path = "$root/$file"
|
||||
val handle = fopen(path, "rb") ?: throw IllegalArgumentException("Resource not found: $path")
|
||||
|
||||
try {
|
||||
if (fseek(handle, 0, SEEK_END) != 0) throw IllegalArgumentException("Resource is not seekable: $path")
|
||||
val size = ftell(handle)
|
||||
if (size < 0L) throw IllegalArgumentException("Cannot determine the size of: $path")
|
||||
if (size == 0L) return ByteArray(0)
|
||||
if (fseek(handle, 0, SEEK_SET) != 0) throw IllegalArgumentException("Cannot rewind: $path")
|
||||
|
||||
val bytes = ByteArray(size.toInt())
|
||||
bytes.usePinned { pinned ->
|
||||
var read = 0
|
||||
while (read < bytes.size) {
|
||||
val count =
|
||||
fread(
|
||||
pinned.addressOf(read),
|
||||
1.convert(),
|
||||
(bytes.size - read).convert(),
|
||||
handle,
|
||||
).toInt()
|
||||
if (count <= 0) break
|
||||
read += count
|
||||
}
|
||||
if (read != bytes.size) {
|
||||
throw IllegalArgumentException("Short read on $path: got $read of ${bytes.size} bytes")
|
||||
}
|
||||
}
|
||||
return bytes
|
||||
} finally {
|
||||
fclose(handle)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Vendored
+140
@@ -0,0 +1,140 @@
|
||||
/*
|
||||
* Copyright (c) 2025 Vitor Pamplona
|
||||
*
|
||||
* Permission is hereby granted, free of charge, to any person obtaining a copy of
|
||||
* this software and associated documentation files (the "Software"), to deal in
|
||||
* the Software without restriction, including without limitation the rights to use,
|
||||
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
|
||||
* Software, and to permit persons to whom the Software is furnished to do so,
|
||||
* subject to the following conditions:
|
||||
*
|
||||
* The above copyright notice and this permission notice shall be included in all
|
||||
* copies or substantial portions of the Software.
|
||||
*
|
||||
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
|
||||
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
|
||||
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
|
||||
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
|
||||
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
|
||||
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
||||
*/
|
||||
package com.vitorpamplona.quartz.utils.cache
|
||||
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFalse
|
||||
import kotlin.test.assertNull
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
/**
|
||||
* Drives every key into a single bucket of the [StripedHashMap] backing this target's
|
||||
* [LargeCache], so the bucket-chain paths run deterministically instead of only when a
|
||||
* hash happens to collide.
|
||||
*
|
||||
* The one that most needs it is removal. Chain nodes hold their `next` immutably — that
|
||||
* is what lets a reader walk a chain with no lock — so removing from the middle has to
|
||||
* clone the nodes ahead of the target onto its tail and publish a new head. With
|
||||
* well-spread keys that path almost never sees a chain longer than two.
|
||||
*/
|
||||
class LargeCacheCollisionTest {
|
||||
/** Every instance lands in the same bucket, and in the same stripe. */
|
||||
private data class Collides(
|
||||
val id: Int,
|
||||
) {
|
||||
override fun hashCode() = 0
|
||||
}
|
||||
|
||||
private fun filled(n: Int) =
|
||||
LargeCache<Collides, Int>().apply {
|
||||
for (i in 0 until n) put(Collides(i), i)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun readsFindEveryEntryInOneChain() {
|
||||
val cache = filled(200)
|
||||
|
||||
assertEquals(200, cache.size())
|
||||
for (i in 0 until 200) {
|
||||
assertEquals(i, cache.get(Collides(i)), "entry $i")
|
||||
assertTrue(cache.containsKey(Collides(i)))
|
||||
}
|
||||
assertNull(cache.get(Collides(200)))
|
||||
assertFalse(cache.containsKey(Collides(200)))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun overwriteInAChainReplacesInPlace() {
|
||||
val cache = filled(200)
|
||||
|
||||
for (i in 0 until 200) cache.put(Collides(i), i * 10)
|
||||
|
||||
assertEquals(200, cache.size(), "overwriting must not lengthen the chain")
|
||||
for (i in 0 until 200) assertEquals(i * 10, cache.get(Collides(i)))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun removeFromTheMiddleKeepsTheRestOfTheChain() {
|
||||
val cache = filled(200)
|
||||
|
||||
// Head, tail and middle of the chain, in an order that leaves gaps behind.
|
||||
for (i in 0 until 200 step 3) {
|
||||
assertEquals(i, cache.remove(Collides(i)), "remove $i returns its value")
|
||||
}
|
||||
|
||||
val expected = (0 until 200).filter { it % 3 != 0 }
|
||||
assertEquals(expected.size, cache.size())
|
||||
for (i in 0 until 200) {
|
||||
if (i % 3 == 0) {
|
||||
assertNull(cache.get(Collides(i)), "entry $i was removed")
|
||||
} else {
|
||||
assertEquals(i, cache.get(Collides(i)), "entry $i survived")
|
||||
}
|
||||
}
|
||||
|
||||
val seen = mutableListOf<Int>()
|
||||
cache.forEach { _, v -> seen.add(v) }
|
||||
assertEquals(expected.toSet(), seen.toSet(), "iteration must match the survivors")
|
||||
assertEquals(expected.size, seen.size, "iteration must not double-count")
|
||||
|
||||
assertNull(cache.remove(Collides(0)), "removing twice is a no-op")
|
||||
assertEquals(expected.size, cache.size())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun getOrCreateAndCreateIfAbsentWalkTheChain() {
|
||||
val cache = filled(200)
|
||||
var builds = 0
|
||||
|
||||
for (i in 0 until 200) {
|
||||
assertEquals(
|
||||
i,
|
||||
cache.getOrCreate(Collides(i)) {
|
||||
builds++
|
||||
-1
|
||||
},
|
||||
)
|
||||
assertFalse(cache.createIfAbsent(Collides(i)) { -1 })
|
||||
}
|
||||
assertEquals(0, builds, "nothing in the chain should have been rebuilt")
|
||||
|
||||
assertTrue(cache.createIfAbsent(Collides(500)) { 500 })
|
||||
assertEquals(500, cache.get(Collides(500)))
|
||||
assertEquals(201, cache.size())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun clearEmptiesAFullBucket() {
|
||||
val cache = filled(200)
|
||||
|
||||
cache.clear()
|
||||
|
||||
assertEquals(0, cache.size())
|
||||
assertTrue(cache.isEmpty())
|
||||
assertNull(cache.get(Collides(7)))
|
||||
assertEquals(0, cache.count { _, _ -> true })
|
||||
|
||||
cache.put(Collides(1), 1)
|
||||
assertEquals(1, cache.size())
|
||||
assertEquals(1, cache.get(Collides(1)))
|
||||
}
|
||||
}
|
||||
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))
|
||||
}
|
||||
}
|
||||
+183
@@ -0,0 +1,183 @@
|
||||
/*
|
||||
* 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.concurrent.atomics.AtomicInt
|
||||
import kotlin.concurrent.atomics.ExperimentalAtomicApi
|
||||
import kotlin.native.concurrent.ObsoleteWorkersApi
|
||||
import kotlin.native.concurrent.TransferMode
|
||||
import kotlin.native.concurrent.Worker
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertNotNull
|
||||
|
||||
/**
|
||||
* Stresses the one path the rest of the suite never reaches: **many writers inserting into
|
||||
* the same bucket at the same time.**
|
||||
*
|
||||
* A striped table is only sound if the stripe is a function of the bucket. Derive the two
|
||||
* from different parts of the hash and the locks stop partitioning the table: two keys can
|
||||
* share a bucket while holding different locks, so two writers read the same chain head and
|
||||
* both publish over it, and one insert vanishes. The first cut of [StripedHashMap] did
|
||||
* exactly that — bucket from the low bits, stripe from bits 16-19 — and nothing here caught
|
||||
* it, because every other suite uses small `Int` keys or a constant `hashCode`, which all
|
||||
* collapse onto stripe 0 and serialise by accident.
|
||||
*
|
||||
* So the keys are built on purpose: groups of four sharing their low 16 bits, which puts a
|
||||
* group in one bucket at any table size this reaches, while differing in bits 16-19 — the
|
||||
* part the broken version striped on. Only [BUCKETS] buckets are used, so chains run to
|
||||
* hundreds of nodes and each insert holds its lock for a while, which is what widens the
|
||||
* window; with well-spread keys it is a few nanoseconds and nothing is observable.
|
||||
*
|
||||
* Being straight about what this is: a stress test, not a deterministic reproducer. It did
|
||||
* not fail against the broken striping in the runs attempted, which says that race is rare
|
||||
* rather than absent — the defect is a plain lost update, provable by reading
|
||||
* [StripedHashMap]'s stripe selection against its bucket index, and was fixed on that
|
||||
* basis. What this test is worth is being the only coverage of concurrent same-bucket
|
||||
* inserts at all, and it would catch a coarser regression.
|
||||
*/
|
||||
@OptIn(ObsoleteWorkersApi::class, ExperimentalAtomicApi::class)
|
||||
class LargeCacheStripingTest {
|
||||
/**
|
||||
* `hashOf` in [StripedHashMap] spreads with `h xor (h ushr 16)`, so to land on a final
|
||||
* hash of `(member shl 16) or bucket` the raw hashCode has to be
|
||||
* `(member shl 16) or (bucket xor member)`. The low 16 bits are then exactly `bucket`,
|
||||
* shared by all four members of a group at every table size this test reaches.
|
||||
*
|
||||
* Only [BUCKETS] distinct buckets are used, so chains run to hundreds of nodes. That
|
||||
* matters: a writer walks its chain *holding the stripe lock*, so a long chain is what
|
||||
* makes the window wide enough for two writers on two different locks to overlap
|
||||
* inside the same bucket. With well-spread keys the window is a few nanoseconds and
|
||||
* the defect hides.
|
||||
*/
|
||||
private data class GroupedKey(
|
||||
val group: Int,
|
||||
val member: Int,
|
||||
) {
|
||||
override fun hashCode(): Int {
|
||||
val bucket = group and (BUCKETS - 1)
|
||||
return (member shl 16) or (bucket xor member)
|
||||
}
|
||||
}
|
||||
|
||||
/** Runs [job] on [WORKERS] threads released together, so they actually contend. */
|
||||
private fun inParallel(job: (workerId: Int) -> Unit) {
|
||||
val ready = AtomicInt(0)
|
||||
val go = AtomicInt(0)
|
||||
val workers = List(WORKERS) { Worker.start() }
|
||||
|
||||
val futures =
|
||||
workers.mapIndexed { id, worker ->
|
||||
worker.execute(TransferMode.SAFE, { Triple(job, id, ready to go) }) { (block, workerId, gates) ->
|
||||
val (readyGate, goGate) = gates
|
||||
readyGate.fetchAndAdd(1)
|
||||
while (goGate.load() == 0) { }
|
||||
block(workerId)
|
||||
}
|
||||
}
|
||||
|
||||
while (ready.load() < WORKERS) { }
|
||||
go.store(1)
|
||||
|
||||
futures.forEach { it.result }
|
||||
workers.forEach { it.requestTermination().result }
|
||||
}
|
||||
|
||||
@Test
|
||||
fun concurrentInsertsAcrossSharedBucketsKeepEveryEntry() {
|
||||
repeat(ROUNDS) { round -> insertRound(round) }
|
||||
}
|
||||
|
||||
private fun insertRound(round: Int) {
|
||||
val cache = LargeCache<GroupedKey, Int>()
|
||||
|
||||
inParallel { worker ->
|
||||
for (group in 0 until GROUPS) {
|
||||
cache.put(GroupedKey(group, worker), group)
|
||||
}
|
||||
}
|
||||
|
||||
val total = GROUPS * WORKERS
|
||||
|
||||
val seen = mutableSetOf<GroupedKey>()
|
||||
cache.forEach { key, _ -> seen.add(key) }
|
||||
|
||||
assertEquals(total, seen.size, "round $round: iteration lost or duplicated entries in a shared bucket")
|
||||
assertEquals(total, cache.size(), "round $round: size() disagrees with what the table holds")
|
||||
|
||||
for (group in 0 until GROUPS) {
|
||||
for (member in 0 until WORKERS) {
|
||||
assertEquals(
|
||||
group,
|
||||
assertNotNull(
|
||||
cache.get(GroupedKey(group, member)),
|
||||
"round $round: entry ($group, $member) was dropped by a concurrent insert",
|
||||
),
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun concurrentCreateIfAbsentAcrossSharedBucketsReportsOneInsertEach() {
|
||||
repeat(ROUNDS) { createIfAbsentRound() }
|
||||
}
|
||||
|
||||
private fun createIfAbsentRound() {
|
||||
val cache = LargeCache<GroupedKey, Int>()
|
||||
val contended = GROUPS / 4
|
||||
|
||||
// Every worker races for the same keys this time, so a lost update shows up as a
|
||||
// duplicate in the chain rather than a missing entry.
|
||||
inParallel { _ ->
|
||||
for (group in 0 until contended) {
|
||||
for (member in 0 until WORKERS) {
|
||||
cache.createIfAbsent(GroupedKey(group, member)) { group }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
val total = contended * WORKERS
|
||||
val seen = mutableSetOf<GroupedKey>()
|
||||
var visited = 0
|
||||
cache.forEach { key, _ ->
|
||||
seen.add(key)
|
||||
visited++
|
||||
}
|
||||
|
||||
assertEquals(total, visited, "a key was inserted twice into the same chain")
|
||||
assertEquals(total, seen.size)
|
||||
assertEquals(total, cache.size())
|
||||
}
|
||||
|
||||
companion object {
|
||||
private const val WORKERS = 4
|
||||
|
||||
/** Large enough that the four workers overlap for essentially the whole run. */
|
||||
private const val GROUPS = 4_096
|
||||
|
||||
/** Few enough that chains grow long and every insert holds its lock for a while. */
|
||||
private const val BUCKETS = 64
|
||||
|
||||
/** Repeated on a fresh table, because only the *insert* path can lose a write. */
|
||||
private const val ROUNDS = 20
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,190 @@
|
||||
/*
|
||||
* 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
|
||||
|
||||
/**
|
||||
* Native actual for [UrlEncoder], shared by linuxX64 and every Apple target.
|
||||
*
|
||||
* ## Why this is not a library call any more
|
||||
*
|
||||
* Both native targets used to delegate to `net.thauvin.erik.urlencoder.UrlEncoderUtil`,
|
||||
* which implements RFC 3986 percent-encoding. The JVM/Android actual is
|
||||
* `java.net.URLEncoder`/`URLDecoder`, which implements
|
||||
* `application/x-www-form-urlencoded`. Those are different specifications, and the
|
||||
* difference was observable in three places:
|
||||
*
|
||||
* ```
|
||||
* JVM/Android UrlEncoderUtil
|
||||
* encode(" ") "+" "%20"
|
||||
* encode("*") "*" "%2A"
|
||||
* decode("a+b") "a b" "a+b"
|
||||
* ```
|
||||
*
|
||||
* That is not cosmetic. [encode] builds strings that leave the device — `TorrentEvent`
|
||||
* puts it in magnet links, `Nip54InlineMetadata` in inline metadata, `Nip47DeepLink` in
|
||||
* the `callback`, `appname` and `value` parameters of NWC deep links — so Android and
|
||||
* iOS were emitting different bytes for the same title. The decode row is worse than
|
||||
* cosmetic: a magnet link or deep link written by Android carries `+` for its spaces,
|
||||
* and reading it on iOS produced a string with literal plus signs instead of spaces, no
|
||||
* error anywhere.
|
||||
*
|
||||
* So this matches `URLEncoder`/`URLDecoder` exactly instead: the unreserved set is
|
||||
* alphanumerics plus `-`, `_`, `.` and `*` (note `*` survives and `~` does not — the
|
||||
* opposite of RFC 3986), space encodes to `+`, everything else to uppercase `%XX` of
|
||||
* its UTF-8 bytes, and decoding maps `+` back to a space. `UrlEncoderTest` in
|
||||
* `commonTest` pins all of it against the JVM on every target.
|
||||
*
|
||||
* Both directions short-circuit the way the `java.net` pair does: a string with nothing
|
||||
* to change is returned as-is rather than rebuilt.
|
||||
*
|
||||
* One deliberate edge difference: an *unpaired* UTF-16 surrogate encodes as `%EF%BF%BD`
|
||||
* (Kotlin's replacement character) where the JVM gives `%3F`. Nostr content is
|
||||
* well-formed UTF-16, and chasing it would cost a scan on every call.
|
||||
*/
|
||||
actual object UrlEncoder {
|
||||
private const val HEX = "0123456789ABCDEF"
|
||||
|
||||
actual fun encode(value: String): String {
|
||||
var index = 0
|
||||
while (index < value.length && isUnreserved(value[index])) index++
|
||||
if (index == value.length) return value
|
||||
|
||||
val result = StringBuilder(value.length + ESCAPE_HEADROOM)
|
||||
result.append(value, 0, index)
|
||||
|
||||
while (index < value.length) {
|
||||
val char = value[index]
|
||||
when {
|
||||
isUnreserved(char) -> {
|
||||
result.append(char)
|
||||
index++
|
||||
}
|
||||
|
||||
char == ' ' -> {
|
||||
result.append('+')
|
||||
index++
|
||||
}
|
||||
|
||||
else -> {
|
||||
// Escaped as a run rather than character by character, so a surrogate
|
||||
// pair becomes one 4-byte sequence instead of two malformed 3-byte ones.
|
||||
val start = index
|
||||
do {
|
||||
index++
|
||||
} while (index < value.length && !isUnreserved(value[index]) && value[index] != ' ')
|
||||
appendEscaped(result, value, start, index)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return result.toString()
|
||||
}
|
||||
|
||||
actual fun decode(value: String): String {
|
||||
if (value.indexOf('%') < 0 && value.indexOf('+') < 0) return value
|
||||
|
||||
val result = StringBuilder(value.length)
|
||||
var index = 0
|
||||
// Sized on first use for the longest run that could still follow, then reused —
|
||||
// one allocation for the whole string, as in URLDecoder.
|
||||
var escaped: ByteArray? = null
|
||||
|
||||
while (index < value.length) {
|
||||
when (val char = value[index]) {
|
||||
'+' -> {
|
||||
result.append(' ')
|
||||
index++
|
||||
}
|
||||
|
||||
'%' -> {
|
||||
val buffer = escaped ?: ByteArray((value.length - index) / 3).also { escaped = it }
|
||||
var count = 0
|
||||
while (index + 2 < value.length && value[index] == '%') {
|
||||
buffer[count++] = decodeEscape(value, index)
|
||||
index += 3
|
||||
}
|
||||
if (index < value.length && value[index] == '%') {
|
||||
throw IllegalArgumentException("URLDecoder: Incomplete trailing escape (%) pattern")
|
||||
}
|
||||
// Decoded as a run so a multi-byte UTF-8 sequence survives.
|
||||
result.append(buffer.decodeToString(0, count))
|
||||
}
|
||||
|
||||
else -> {
|
||||
result.append(char)
|
||||
index++
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return result.toString()
|
||||
}
|
||||
|
||||
/** `URLEncoder`'s `dontNeedEncoding` set: alphanumerics plus these four, and only these. */
|
||||
private fun isUnreserved(char: Char): Boolean =
|
||||
char in 'a'..'z' ||
|
||||
char in 'A'..'Z' ||
|
||||
char in '0'..'9' ||
|
||||
char == '-' ||
|
||||
char == '_' ||
|
||||
char == '.' ||
|
||||
char == '*'
|
||||
|
||||
private fun appendEscaped(
|
||||
result: StringBuilder,
|
||||
value: String,
|
||||
start: Int,
|
||||
end: Int,
|
||||
) {
|
||||
val bytes = value.substring(start, end).encodeToByteArray()
|
||||
for (byte in bytes) {
|
||||
val code = byte.toInt()
|
||||
result.append('%')
|
||||
result.append(HEX[(code shr 4) and 0xF])
|
||||
result.append(HEX[code and 0xF])
|
||||
}
|
||||
}
|
||||
|
||||
private fun decodeEscape(
|
||||
value: String,
|
||||
index: Int,
|
||||
): Byte {
|
||||
val high = hexDigit(value[index + 1])
|
||||
val low = hexDigit(value[index + 2])
|
||||
if (high < 0 || low < 0) {
|
||||
throw IllegalArgumentException(
|
||||
"URLDecoder: Illegal hex characters in escape (%) pattern - ${value.substring(index, index + 3)}",
|
||||
)
|
||||
}
|
||||
return ((high shl 4) or low).toByte()
|
||||
}
|
||||
|
||||
private fun hexDigit(char: Char): Int =
|
||||
when (char) {
|
||||
in '0'..'9' -> char - '0'
|
||||
in 'a'..'f' -> char - 'a' + 10
|
||||
in 'A'..'F' -> char - 'A' + 10
|
||||
else -> -1
|
||||
}
|
||||
|
||||
/** Enough for a handful of escapes before the builder has to grow. */
|
||||
private const val ESCAPE_HEADROOM = 16
|
||||
}
|
||||
+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