mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
perf(commons): halve the list observables' write cost, and measure it
The copy-on-write rewrite of the two list observables was justified on portability -- `java.util.concurrent` has no KMP equivalent, which is what kept them in `jvmAndroid` behind an iOS stub. That says nothing about speed, and these are hot: every consumed event is offered to each observer whose filter could match it, from each relay's socket coroutine. So measure it. `ObserverListBenchmark` keeps the `ConcurrentSkipListSet` implementation verbatim as the baseline and as a differential oracle -- `bothImplementationsAgree` asserts the two emit identical lists for identical input at limits null/50/400, which is the property the rewrite had to preserve. Runs in ~7s, asserts on correctness only, never on wall time. The first run said the rewrite was a regression: slower on the limited insert path and 4.9x slower on 8 threads, with no thread scaling. The cause was one line, not the design. `State.plus` published `ids + added`, and over the limit `ids + added - dropped`; each operator allocates a full copy of the set, so the eviction path rebuilt membership twice per insert. One `HashSet`, sized up front and mutated before publishing, is the fix. After it, against the implementation it replaced: re-delivering an already listed note 3.3-4.1x faster, insert into a populated unlimited list (what nearly every observer registers) 1.3-1.5x faster and widening with n, cold fill at n=5000 1.8x faster, concurrent inserts at parity on one thread and 2.1-2.7x slower on 4-8. That last number is a real cost and is written on both classes rather than left out: threads serialize on one reference and a lost CAS discards its copy, where the skip list striped across keys and scaled. Two things bound it -- the benchmark's threads do nothing but insert, while real ingest spends most of its per-event budget verifying signatures and parsing before it reaches an observer; and nearly every production observer registers with no limit, which is the shape copy-on-write wins. Also recorded, since it framed the earlier discussion wrong: lock-free was never the difference between the two. The skip-list version was lock-free too -- ConcurrentHashMap stripes per key, so `compute` holds a bin, not a monitor. The choice was portability and speed. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01KkULS5SVq4GHDdoCzajKi8
This commit is contained in:
@@ -1219,3 +1219,86 @@ Verified: `:commons:compileCommonMainKotlinMetadata`,
|
||||
`:commons:jvmTest`, `:desktopApp:test`, `:cli:test`,
|
||||
`:amethyst:compileFdroidDebugKotlin`, `:amethyst:testPlayDebugUnitTest`,
|
||||
`spotlessCheck`.
|
||||
|
||||
## Step 8 — measuring the copy-on-write observables (the "is it actually faster?" question)
|
||||
|
||||
The rewrite of `NoteListMatchingFilter` / `EventListMatchingFilter` from
|
||||
`ConcurrentSkipListSet` + `ConcurrentHashMap` to copy-on-write over an
|
||||
`AtomicReference` was made for **portability** — `java.util.concurrent` has no
|
||||
KMP equivalent, and that is what kept these two in `jvmAndroid` with an iOS
|
||||
stub. That argument says nothing about speed, and these sit on a hot path:
|
||||
every consumed event is offered to each observer whose filter could match it,
|
||||
from each relay's socket coroutine. A few hundred events a second across a
|
||||
dozen relays reaches this code thousands of times a second.
|
||||
|
||||
So it was measured rather than argued. `ObserverListBenchmark`
|
||||
(`commons/src/jvmTest/.../prodbench/`) keeps the skip-list implementation
|
||||
verbatim as the baseline and as a **differential oracle**: `bothImplementations
|
||||
Agree` asserts the two emit identical lists for identical input at limits
|
||||
`null` / 50 / 400, which is the property the rewrite had to preserve. It runs
|
||||
in ~7s and asserts only on correctness, never on wall time.
|
||||
|
||||
**What the first run found:** the copy-on-write version was *slower* on the
|
||||
limited-filter insert path and **4.9x slower on 8 concurrent threads**, with
|
||||
zero thread scaling. The cause was not the design but one line —
|
||||
`State.plus` published `ids + added` and, over the limit, `ids + added -
|
||||
dropped`. Each operator allocates a full copy of the set, so the eviction path
|
||||
rebuilt the membership set **twice per insert**. Replacing both with a single
|
||||
`HashSet` copy (sized up front, mutated, then published) is the whole of the
|
||||
fix.
|
||||
|
||||
**After that fix**, against the implementation it replaced:
|
||||
|
||||
| shape | result |
|
||||
|---|---|
|
||||
| re-deliver an already listed note (steady state once a screen is warm) | **3.3–4.1x faster** |
|
||||
| insert into a populated *unlimited* list, n = 100 / 1000 | **1.3–1.5x faster**, widening with n |
|
||||
| cold fill, n = 5000 | **1.8x faster** |
|
||||
| concurrent inserts into the same observer, 1 thread | parity |
|
||||
| concurrent inserts into the same observer, 4 / 8 threads | **2.1–2.7x slower** |
|
||||
|
||||
The last row is the real cost and is documented on both classes rather than
|
||||
buried here: threads serialize on one reference and a lost CAS discards its
|
||||
copy, where the skip list striped across keys and scaled with thread count. Two
|
||||
things bound it. The benchmark's threads do nothing but insert, while real
|
||||
ingest spends most of its per-event budget on signature verification and
|
||||
parsing before reaching an observer, so the contention window is a fraction of
|
||||
the measured one. And nearly every production observer registers with **no
|
||||
`limit`** — grep the `observeNotes` / `observeEvents` call sites — which is the
|
||||
shape copy-on-write wins.
|
||||
|
||||
Worth revisiting if a profile ever disagrees. The lever would be the emit, not
|
||||
the lock: both implementations already materialize the whole list on every
|
||||
write, so an observer that only ever appends is paying O(n) to tell the UI
|
||||
about one new row.
|
||||
|
||||
**Also worth recording:** "lock-free" was never the differentiator between the
|
||||
two. The skip-list version was lock-free too — `ConcurrentHashMap` stripes per
|
||||
key, so `compute` holds one bin, not a monitor. The choice was portability and
|
||||
speed, not locking.
|
||||
|
||||
Verified: `:commons:jvmTest`, `:commons:compileIosMainKotlinMetadata`,
|
||||
`:commons:spotlessApply`.
|
||||
|
||||
## Step 9 — the weak note cache has no strong referent in tests
|
||||
|
||||
`compose-ui-test` went red on `DesktopCachePipelineTest`:
|
||||
`FollowingFeedFilter only includes notes from followed users`, `expected:<1> but
|
||||
was:<0>`, and it would not reproduce on a dev machine.
|
||||
|
||||
The cause is this branch, indirectly. `DesktopLocalCache` kept a
|
||||
`notesByAuthor: ConcurrentHashMap<HexKey, MutableSet<Note>>` index for metadata
|
||||
invalidation — an unbounded strong map holding every note the Desktop app ever
|
||||
saw, which quietly defeated the point of storing them in a `LargeSoftCache`.
|
||||
Removing it was right (the Android cache never had one). But it was also the
|
||||
only thing keeping the test fixtures alive: every test consumes events and then
|
||||
queries the cache for the notes they produced, with nothing in between holding a
|
||||
reference. A GC landing in that window empties the cache.
|
||||
|
||||
It reproduces deterministically with two `System.gc()` calls before the query,
|
||||
and it is not specific to that one test — all 46 consume sites have the shape,
|
||||
so CI's tighter heap just picked the victim. The fix routes every consume
|
||||
through an `ingest` helper that pins what the cache built for the lifetime of
|
||||
the test instance, the way a screen holds the notes it is showing in the app,
|
||||
and keeps the forced GC in the test that failed so the contract is asserted
|
||||
rather than left to the heap.
|
||||
|
||||
+26
-2
@@ -56,6 +56,21 @@ import kotlin.concurrent.atomics.ExperimentalAtomicApi
|
||||
* concurrent structure plus a `HashSet` to strip the duplicates that walk could surface.
|
||||
*
|
||||
* The copy is O(n) per write, which is what the previous per-emit snapshot already cost.
|
||||
*
|
||||
* **Measured**, against the `ConcurrentSkipListSet` implementation this replaced — see
|
||||
* `ObserverListBenchmark`, which keeps that implementation as the baseline and as a
|
||||
* differential oracle:
|
||||
* - re-delivery of an already listed note (the steady state once a screen is warm): 3.3–4.1x
|
||||
* faster, because it is a load plus a set lookup where the skip list took a `compute`;
|
||||
* - insert into a populated unlimited list (what nearly every observer registers): 1.3–1.5x
|
||||
* faster at n = 100–1000, widening with n;
|
||||
* - concurrent inserts into the SAME observer: at parity on one thread, 2.1–2.7x slower on
|
||||
* 4–8. That is the honest cost of this design: threads serialize on one reference and a
|
||||
* lost CAS throws away its copy, where the skip list striped across keys and scaled. It is
|
||||
* measured with threads doing nothing but inserting; real ingest spends most of its
|
||||
* per-event budget verifying signatures and parsing before it reaches an observer, so the
|
||||
* contention window is a fraction of that. Worth revisiting if a profile ever says
|
||||
* otherwise.
|
||||
*/
|
||||
class EventListMatchingFilter<T : Event>(
|
||||
private val filter: Filter,
|
||||
@@ -106,12 +121,21 @@ class EventListMatchingFilter<T : Event>(
|
||||
grown.add(entry)
|
||||
grown.addAll(entries.subList(at, entries.size))
|
||||
|
||||
if (limit == null || grown.size <= limit) return State(grown, ids + entry.note.idHex)
|
||||
// One copy of the membership set per write. `ids + added` and `ids + added - dropped`
|
||||
// read well but allocate one full set per operator, so the limit path paid for two;
|
||||
// measured, that was the whole of this implementation's deficit against the skip list
|
||||
// on a limited filter.
|
||||
val nextIds = HashSet<HexKey>(((entries.size + 2) / 0.75f).toInt() + 1)
|
||||
nextIds.addAll(ids)
|
||||
nextIds.add(entry.note.idHex)
|
||||
|
||||
if (limit == null || grown.size <= limit) return State(grown, nextIds)
|
||||
|
||||
// Over the limit: drop the oldest, which sorts last. That can be the entry just
|
||||
// inserted, and then it is simply not listed — as the previous pollLast() did.
|
||||
val dropped = grown.removeAt(grown.size - 1)
|
||||
return State(grown, ids + entry.note.idHex - dropped.note.idHex)
|
||||
nextIds.remove(dropped.note.idHex)
|
||||
return State(grown, nextIds)
|
||||
}
|
||||
|
||||
private fun State.minus(idHex: HexKey): State = State(entries.filterNot { it.note.idHex == idHex }, ids - idHex)
|
||||
|
||||
+26
-2
@@ -56,6 +56,21 @@ import kotlin.concurrent.atomics.ExperimentalAtomicApi
|
||||
* concurrent structure plus a `HashSet` to strip the duplicates that walk could surface.
|
||||
*
|
||||
* The copy is O(n) per write, which is what the previous per-emit snapshot already cost.
|
||||
*
|
||||
* **Measured**, against the `ConcurrentSkipListSet` implementation this replaced — see
|
||||
* `ObserverListBenchmark`, which keeps that implementation as the baseline and as a
|
||||
* differential oracle:
|
||||
* - re-delivery of an already listed note (the steady state once a screen is warm): 3.3–4.1x
|
||||
* faster, because it is a load plus a set lookup where the skip list took a `compute`;
|
||||
* - insert into a populated unlimited list (what nearly every observer registers): 1.3–1.5x
|
||||
* faster at n = 100–1000, widening with n;
|
||||
* - concurrent inserts into the SAME observer: at parity on one thread, 2.1–2.7x slower on
|
||||
* 4–8. That is the honest cost of this design: threads serialize on one reference and a
|
||||
* lost CAS throws away its copy, where the skip list striped across keys and scaled. It is
|
||||
* measured with threads doing nothing but inserting; real ingest spends most of its
|
||||
* per-event budget verifying signatures and parsing before it reaches an observer, so the
|
||||
* contention window is a fraction of that. Worth revisiting if a profile ever says
|
||||
* otherwise.
|
||||
*/
|
||||
class NoteListMatchingFilter(
|
||||
private val filter: Filter,
|
||||
@@ -109,12 +124,21 @@ class NoteListMatchingFilter(
|
||||
grown.add(entry)
|
||||
grown.addAll(entries.subList(at, entries.size))
|
||||
|
||||
if (limit == null || grown.size <= limit) return State(grown, ids + entry.note.idHex)
|
||||
// One copy of the membership set per write. `ids + added` and `ids + added - dropped`
|
||||
// read well but allocate one full set per operator, so the limit path paid for two;
|
||||
// measured, that was the whole of this implementation's deficit against the skip list
|
||||
// on a limited filter.
|
||||
val nextIds = HashSet<HexKey>(((entries.size + 2) / 0.75f).toInt() + 1)
|
||||
nextIds.addAll(ids)
|
||||
nextIds.add(entry.note.idHex)
|
||||
|
||||
if (limit == null || grown.size <= limit) return State(grown, nextIds)
|
||||
|
||||
// Over the limit: drop the oldest, which sorts last. That can be the entry just
|
||||
// inserted, and then it is simply not listed — as the previous pollLast() did.
|
||||
val dropped = grown.removeAt(grown.size - 1)
|
||||
return State(grown, ids + entry.note.idHex - dropped.note.idHex)
|
||||
nextIds.remove(dropped.note.idHex)
|
||||
return State(grown, nextIds)
|
||||
}
|
||||
|
||||
private fun State.minus(idHex: HexKey): State = State(entries.filterNot { it.note.idHex == idHex }, ids - idHex)
|
||||
|
||||
+377
@@ -0,0 +1,377 @@
|
||||
/*
|
||||
* 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.amethyst.commons.prodbench
|
||||
|
||||
import com.vitorpamplona.amethyst.commons.model.AddressableNote
|
||||
import com.vitorpamplona.amethyst.commons.model.Note
|
||||
import com.vitorpamplona.amethyst.commons.model.observables.NoteListMatchingFilter
|
||||
import com.vitorpamplona.amethyst.commons.model.observables.Observable
|
||||
import com.vitorpamplona.quartz.nip01Core.core.AddressableEvent
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip10Notes.TextNoteEvent
|
||||
import com.vitorpamplona.quartz.utils.EventFactory
|
||||
import java.util.Random
|
||||
import java.util.concurrent.ConcurrentHashMap
|
||||
import java.util.concurrent.ConcurrentSkipListSet
|
||||
import java.util.concurrent.CountDownLatch
|
||||
import java.util.concurrent.atomic.AtomicLong
|
||||
import kotlin.concurrent.thread
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
|
||||
/**
|
||||
* Measures [NoteListMatchingFilter], the list every `observeNotes` / `observeEvents` screen
|
||||
* sits behind, against the `ConcurrentSkipListSet` implementation it replaced.
|
||||
*
|
||||
* It is on a hot path: every consumed event is offered to each observer whose filter could
|
||||
* match it, from each relay's socket coroutine, so a few hundred events a second across a
|
||||
* dozen relays reach this code thousands of times a second. The rewrite to copy-on-write was
|
||||
* made to get the class into `commonMain` (`java.util.concurrent` has no KMP equivalent);
|
||||
* that is a portability argument, and this file is the performance one, which is a separate
|
||||
* question and has to be answered with numbers rather than asymptotics.
|
||||
*
|
||||
* [SkipListReference] is the previous implementation, kept verbatim as the baseline. It also
|
||||
* serves as a differential oracle: [bothImplementationsAgree] asserts the two produce identical
|
||||
* output for identical input, which is the property the rewrite actually had to preserve.
|
||||
*
|
||||
* Deterministic and offline. Prints ns/op; asserts only on correctness, never on wall time,
|
||||
* since CI machines vary.
|
||||
*/
|
||||
class ObserverListBenchmark {
|
||||
private val author = "d0d0a746b44c9de8422165aef520b1fe041eedf5794f7592505477eeac122c18"
|
||||
|
||||
/** Touches the emitted list so neither implementation's O(n) materialization is dead code. */
|
||||
private val blackhole = AtomicLong()
|
||||
|
||||
private fun sink(list: List<Note>) {
|
||||
blackhole.addAndGet(list.size.toLong() + if (list.isEmpty()) 0 else list[0].idHex.length.toLong())
|
||||
}
|
||||
|
||||
private fun event(
|
||||
i: Int,
|
||||
createdAt: Long,
|
||||
): Event =
|
||||
EventFactory.create(
|
||||
id = "%064x".format(i),
|
||||
pubKey = author,
|
||||
createdAt = createdAt,
|
||||
kind = TextNoteEvent.KIND,
|
||||
tags = emptyArray(),
|
||||
content = "n",
|
||||
sig = "00".repeat(64),
|
||||
)
|
||||
|
||||
/**
|
||||
* createdAt is deliberately NOT monotonic in i: a real feed arrives out of order, so the
|
||||
* insert lands mid-list rather than always at the head, which is the cheapest case for
|
||||
* both implementations and would flatter whichever one is worse.
|
||||
*/
|
||||
private fun fixtures(n: Int): List<Pair<Event, Note>> {
|
||||
val rnd = Random(42)
|
||||
return (0 until n).map { i ->
|
||||
val e = event(i, 1_700_000_000L + rnd.nextInt(n * 4))
|
||||
val note = Note(e.id)
|
||||
note.event = e
|
||||
e to note
|
||||
}
|
||||
}
|
||||
|
||||
private fun cow(f: Filter) = NoteListMatchingFilter(f, { emptyList() }, ::sink)
|
||||
|
||||
private fun skipList(f: Filter) = SkipListReference(f, ::sink)
|
||||
|
||||
// ---------------------------------------------------------------------- harness
|
||||
|
||||
private fun bench(
|
||||
label: String,
|
||||
reps: Int = 5,
|
||||
warmups: Int = 2,
|
||||
body: () -> Int,
|
||||
): Double {
|
||||
repeat(warmups) { body() }
|
||||
val times = ArrayList<Double>(reps)
|
||||
repeat(reps) {
|
||||
val start = System.nanoTime()
|
||||
val ops = body()
|
||||
times.add((System.nanoTime() - start).toDouble() / ops)
|
||||
}
|
||||
times.sort()
|
||||
val median = times[times.size / 2]
|
||||
println(" %-44s %9.0f ns/op %11.0f ops/s".format(label, median, 1_000_000_000.0 / median))
|
||||
return median
|
||||
}
|
||||
|
||||
private fun compare(
|
||||
name: String,
|
||||
cowNs: Double,
|
||||
skipNs: Double,
|
||||
) {
|
||||
val ratio = cowNs / skipNs
|
||||
val verdict =
|
||||
when {
|
||||
ratio < 0.95 -> "copy-on-write faster x%.2f".format(1 / ratio)
|
||||
ratio > 1.05 -> "copy-on-write slower x%.2f".format(ratio)
|
||||
else -> "parity"
|
||||
}
|
||||
println(" -> $name: $verdict\n")
|
||||
}
|
||||
|
||||
// ------------------------------------------------------------------ correctness
|
||||
|
||||
@Test
|
||||
fun bothImplementationsAgree() {
|
||||
val fx = fixtures(500)
|
||||
for (limit in listOf(null, 50, 400)) {
|
||||
val f = Filter(kinds = listOf(TextNoteEvent.KIND), limit = limit)
|
||||
|
||||
var fromCow: List<Note> = emptyList()
|
||||
val a = NoteListMatchingFilter(f, { emptyList() }, { fromCow = it })
|
||||
var fromSkipList: List<Note> = emptyList()
|
||||
val b = SkipListReference(f) { fromSkipList = it }
|
||||
|
||||
fx.forEach { (e, note) -> a.new(e, note) }
|
||||
fx.forEach { (e, note) -> b.new(e, note) }
|
||||
|
||||
assertEquals(
|
||||
fromSkipList.map { it.idHex },
|
||||
fromCow.map { it.idHex },
|
||||
"copy-on-write must emit exactly what the skip list emitted (limit=$limit)",
|
||||
)
|
||||
assertEquals(fromCow.size, fromCow.map { it.idHex }.toSet().size, "no duplicate keys")
|
||||
|
||||
// Re-delivering an already listed note must be a no-op in both.
|
||||
fx.take(50).forEach { (e, note) -> a.new(e, note) }
|
||||
fx.take(50).forEach { (e, note) -> b.new(e, note) }
|
||||
assertEquals(fromSkipList.map { it.idHex }, fromCow.map { it.idHex }, "re-delivery must not change either list")
|
||||
}
|
||||
}
|
||||
|
||||
// ------------------------------------------------------------------- serial load
|
||||
|
||||
@Test
|
||||
fun insertIntoAPopulatedList() {
|
||||
// No limit: that is what nearly every production observer passes, so the list grows
|
||||
// to whatever the cache holds for that kind.
|
||||
println("\n== insert into a list already holding n (seed excluded from timing, no limit) ==")
|
||||
for (n in listOf(100, 1000)) {
|
||||
val seed = fixtures(n)
|
||||
val arrivals = fixtures(n * 2).drop(n).map { (e, _) -> e }
|
||||
val f = Filter(kinds = listOf(TextNoteEvent.KIND))
|
||||
|
||||
fun run(make: () -> Observable): () -> Int =
|
||||
{
|
||||
val subject = make()
|
||||
seed.forEach { (e, note) -> subject.new(e, note) }
|
||||
val notes =
|
||||
arrivals.map { e ->
|
||||
val note = Note(e.id)
|
||||
note.event = e
|
||||
e to note
|
||||
}
|
||||
val t0 = System.nanoTime()
|
||||
notes.forEach { (e, note) -> subject.new(e, note) }
|
||||
elapsed.addAndGet(System.nanoTime() - t0)
|
||||
arrivals.size
|
||||
}
|
||||
|
||||
val c = benchInner("copy-on-write insert @$n", run { cow(f) })
|
||||
val s = benchInner("skip list insert @$n", run { skipList(f) })
|
||||
compare("insert @n=$n", c, s)
|
||||
}
|
||||
}
|
||||
|
||||
private val elapsed = AtomicLong()
|
||||
|
||||
/** Reports only the span the body accumulated into [elapsed], excluding its setup. */
|
||||
private fun benchInner(
|
||||
label: String,
|
||||
body: () -> Int,
|
||||
reps: Int = 5,
|
||||
warmups: Int = 2,
|
||||
): Double {
|
||||
repeat(warmups) {
|
||||
elapsed.set(0)
|
||||
body()
|
||||
}
|
||||
val times = ArrayList<Double>(reps)
|
||||
repeat(reps) {
|
||||
elapsed.set(0)
|
||||
val ops = body()
|
||||
times.add(elapsed.get().toDouble() / ops)
|
||||
}
|
||||
times.sort()
|
||||
val median = times[times.size / 2]
|
||||
println(" %-44s %9.0f ns/op %11.0f ops/s".format(label, median, 1_000_000_000.0 / median))
|
||||
return median
|
||||
}
|
||||
|
||||
@Test
|
||||
fun reDeliveryOfAnAlreadyListedNote() {
|
||||
// The dominant steady-state call once a screen is warm: a newer version of an
|
||||
// addressable, or the same note arriving from a second relay.
|
||||
println("\n== re-delivery of a note already in the list ==")
|
||||
for (n in listOf(100, 1000)) {
|
||||
val fx = fixtures(n)
|
||||
val f = Filter(kinds = listOf(TextNoteEvent.KIND), limit = n)
|
||||
val reps = 100_000
|
||||
|
||||
val a = cow(f).also { s -> fx.forEach { (e, note) -> s.new(e, note) } }
|
||||
val b = skipList(f).also { s -> fx.forEach { (e, note) -> s.new(e, note) } }
|
||||
|
||||
val c =
|
||||
bench("copy-on-write re-deliver @$n") {
|
||||
for (i in 0 until reps) {
|
||||
val (e, note) = fx[i % n]
|
||||
a.new(e, note)
|
||||
}
|
||||
reps
|
||||
}
|
||||
val s =
|
||||
bench("skip list re-deliver @$n") {
|
||||
for (i in 0 until reps) {
|
||||
val (e, note) = fx[i % n]
|
||||
b.new(e, note)
|
||||
}
|
||||
reps
|
||||
}
|
||||
compare("re-deliver n=$n", c, s)
|
||||
}
|
||||
}
|
||||
|
||||
// --------------------------------------------------------------- concurrent load
|
||||
|
||||
private fun concurrently(
|
||||
threads: Int,
|
||||
ops: Int,
|
||||
subject: Observable,
|
||||
fx: List<Pair<Event, Note>>,
|
||||
): Int {
|
||||
val start = CountDownLatch(1)
|
||||
val done = CountDownLatch(threads)
|
||||
repeat(threads) { t ->
|
||||
thread {
|
||||
start.await()
|
||||
var i = t
|
||||
while (i < ops) {
|
||||
val (e, note) = fx[i % fx.size]
|
||||
subject.new(e, note)
|
||||
i += threads
|
||||
}
|
||||
done.countDown()
|
||||
}
|
||||
}
|
||||
start.countDown()
|
||||
done.await()
|
||||
return ops
|
||||
}
|
||||
|
||||
@Test
|
||||
fun concurrentInsertsFromSeveralRelayThreads() {
|
||||
// Worst case on purpose: these threads do nothing but insert. Real ingest spends most
|
||||
// of its per-event budget on signature verification and parsing before reaching an
|
||||
// observer, so actual contention on one filter is a fraction of this.
|
||||
println("\n== concurrent inserts, threads doing nothing but inserting (worst case) ==")
|
||||
val ops = 4_000
|
||||
val fx = fixtures(ops)
|
||||
val f = Filter(kinds = listOf(TextNoteEvent.KIND), limit = 1000)
|
||||
for (threads in listOf(1, 4, 8)) {
|
||||
val c = bench("copy-on-write concurrent x$threads", reps = 3, warmups = 1) { concurrently(threads, ops, cow(f), fx) }
|
||||
val s = bench("skip list concurrent x$threads", reps = 3, warmups = 1) { concurrently(threads, ops, skipList(f), fx) }
|
||||
compare("concurrent x$threads", c, s)
|
||||
}
|
||||
println(" blackhole=${blackhole.get()}")
|
||||
}
|
||||
|
||||
/**
|
||||
* [NoteListMatchingFilter] as it stood before the copy-on-write rewrite, verbatim apart from
|
||||
* dropping the unused `atOnce`/`init` pair. Kept as the baseline these numbers are measured
|
||||
* against, and as the oracle in [bothImplementationsAgree].
|
||||
*
|
||||
* It is lock-free too -- `byId` is a ConcurrentHashMap, striped per key, and every write to
|
||||
* `sorted` for an idHex happens inside that key's `compute` section -- which is why the
|
||||
* comparison is about speed and portability, not about locking.
|
||||
*/
|
||||
private class SkipListReference(
|
||||
private val filter: Filter,
|
||||
private val update: (List<Note>) -> Unit,
|
||||
) : Observable {
|
||||
private class Entry(
|
||||
val note: Note,
|
||||
val createdAt: Long,
|
||||
val id: HexKey,
|
||||
)
|
||||
|
||||
private val order =
|
||||
Comparator<Entry> { a, b ->
|
||||
val byCreatedAt = b.createdAt.compareTo(a.createdAt)
|
||||
if (byCreatedAt != 0) byCreatedAt else a.id.compareTo(b.id)
|
||||
}
|
||||
|
||||
private val sorted = ConcurrentSkipListSet(order)
|
||||
private val byId = ConcurrentHashMap<HexKey, Entry>()
|
||||
|
||||
private fun entryFor(note: Note) = Entry(note, note.createdAt() ?: Long.MIN_VALUE, note.event?.id ?: note.idHex)
|
||||
|
||||
override fun new(
|
||||
event: Event,
|
||||
note: Note,
|
||||
) {
|
||||
if (event is AddressableEvent && note !is AddressableNote) return
|
||||
if (!filter.match(event)) return
|
||||
|
||||
var added = false
|
||||
byId.compute(note.idHex) { _, existing ->
|
||||
existing ?: entryFor(note).also {
|
||||
sorted.add(it)
|
||||
added = true
|
||||
}
|
||||
}
|
||||
if (!added) return
|
||||
|
||||
val limit = filter.limit
|
||||
if (limit != null && sorted.size > limit) {
|
||||
sorted.pollLast()?.let { byId.remove(it.note.idHex, it) }
|
||||
}
|
||||
|
||||
update(snapshot())
|
||||
}
|
||||
|
||||
override fun remove(note: Note) {
|
||||
var removed = false
|
||||
byId.compute(note.idHex) { _, existing ->
|
||||
if (existing != null) {
|
||||
sorted.remove(existing)
|
||||
removed = true
|
||||
}
|
||||
null
|
||||
}
|
||||
if (removed) update(snapshot())
|
||||
}
|
||||
|
||||
/** The skip list's iterator is only weakly consistent, so this had to strip duplicates. */
|
||||
private fun snapshot(): List<Note> {
|
||||
val seen = HashSet<HexKey>()
|
||||
return sorted.mapNotNull { e -> e.note.takeIf { seen.add(it.idHex) } }
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user