mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 03:38:23 +00:00
fix(quartz): stripe from the bucket, not from unrelated hash bits
Audit finding, and a real defect in the striped table two commits back. A striped hash table is only sound when the stripe is a function of the bucket. lockFor picked bits 16-19 of the hash while the bucket index used the low bits, so the 16 locks did not partition the table: two keys could share a bucket while holding different locks, and two writers would then read the same chain head and both publish over it. One insert silently disappears while entryCount counts both. The same window loses an overwrite, and loses entries through remove's chain rebuild. That is precisely the class of bug this work set out to remove from the copy-on-write version it replaced. Stripe now comes from `hash and (STRIPES - 1)`. Because STRIPES and every capacity are powers of two with STRIPES <= capacity, those are exactly the low bits of the bucket index, so same bucket implies same stripe at every size. It stays derived from the hash rather than the capacity, so a key keeps its stripe across a resize, which is what lets growTable exclude writers by taking all of them. INITIAL_CAPACITY is now defined as STRIPES so raising one cannot silently break the invariant. That definition also fixes a memory regression the audit caught: the table allocated 1024 slots eagerly, about 8 KB, per instance. LargeCache is not only the one big LocalCache — EphemeralRoom, RelaySession, PoolRequests and others build one per room, per connection and per subscription set, so a client holds hundreds that stay nearly empty. An empty instance goes from ~8 KB to ~970 bytes. Growth is geometric, so a table that does fill to 100k pays the same ~2n node rebuilds either way; re-measuring the shipped code confirms it (fill 16ms, overwrite 4ms, reads 1ms, 20 scans 16ms, mixed 86ms, 1 GC — unchanged within noise). The KDoc table is updated to those numbers. Adds LargeCacheStripingTest, which builds keys that share a bucket while differing in bits 16-19 and drives four workers at them behind a start barrier, with few enough buckets that chains grow long and each insert holds its lock for a while. It is documented for what it is: a stress test of the concurrent same-bucket path, not a deterministic reproducer — it did not fail against the broken striping in the runs attempted, which makes that race rare rather than absent. The fix rests on reading the stripe selection against the bucket index, not on a red test. Remaining known cost, noted in the KDoc rather than changed here: those ~970 bytes are nearly all the 16 PlatformLocks, two objects each. Folding them into one AtomicIntArray would reach ~250 bytes, but hand-rolling the spin wants its own review rather than a change on the way to merge. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01HxQ1QuyzSkR38iFHbREjoS
This commit is contained in:
+41
-8
@@ -58,14 +58,23 @@ import kotlin.concurrent.atomics.ExperimentalAtomicApi
|
||||
* 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 15 3 3 13 64 1 43MB
|
||||
* 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.
|
||||
* 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
|
||||
*
|
||||
@@ -111,10 +120,23 @@ internal class StripedHashMap<K, V> {
|
||||
}
|
||||
|
||||
/**
|
||||
* Derived from the hash alone, never from the table size, so a key keeps the same
|
||||
* stripe across a resize.
|
||||
* **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 ushr 16) and (STRIPES - 1)]
|
||||
private fun lockFor(hash: Int) = locks[hash and (STRIPES - 1)]
|
||||
|
||||
fun size(): Int = entryCount.load()
|
||||
|
||||
@@ -301,8 +323,19 @@ internal class StripedHashMap<K, V> {
|
||||
*/
|
||||
private const val STRIPES = 16
|
||||
|
||||
/** Sized to carry a warm cache's first few thousand entries without a resize. */
|
||||
private const val INITIAL_CAPACITY = 1024
|
||||
/**
|
||||
* 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
|
||||
}
|
||||
|
||||
+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
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user