refactor(commons): make the list observables lock-free and common

Both matching-filter observables were a ConcurrentSkipListSet coordinated with
a ConcurrentHashMap.compute per-key critical section — lock-free, and with no
Kotlin/Native equivalent, which is what forced the expect/actual and the two
throwing iOS stubs. Copy-on-write over an atomic reference gives the same
lock-freedom with none of that: one immutable State behind an AtomicReference,
each mutation building the next state and publishing it with a compare-and-set,
retrying if another thread won. It is the shape quartz already uses for its
native ConcurrentMap actual, and it compiles everywhere, so the expect/actual,
the jvmAndroid actuals and the iOS stubs are all gone.

Two hand-written workarounds go with them. The emitted snapshot no longer
de-duplicates by idHex — that pass existed only because a weakly-consistent
skip-list iterator can transiently surface a key twice, and a reader of an
immutable state cannot. And emitting is now one map over the published state
rather than a walk of a concurrent structure plus a HashSet, so the per-event
cost drops even though the copy is O(n): the previous code already paid O(n)
per emit, on top of the insert.

The sort key is still snapshotted into an immutable Entry at insertion, so a
replaceable event mutating created_at on the same AddressableNote instance
still cannot move an entry — the test that pins that behaviour, and both 8
thread x 5000 iteration concurrency tests, pass unchanged.

ensureMintDirectoryBackfilled loses its lock too: a CAS on the claim flag means
exactly one caller sweeps, and a second returns immediately instead of blocking
for the length of a full cache scan. It may see a partially filled index, which
is what it already saw before any backfill ran.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01KkULS5SVq4GHDdoCzajKi8
This commit is contained in:
Claude
2026-09-21 17:57:14 +00:00
parent 61ceaa728f
commit c25cf6524f
7 changed files with 285 additions and 469 deletions
@@ -58,8 +58,6 @@ import com.vitorpamplona.amethyst.commons.model.privateChats.ChatroomList
import com.vitorpamplona.amethyst.commons.model.redirectStrayRelayGroupContent
import com.vitorpamplona.amethyst.commons.service.BundledInsert
import com.vitorpamplona.amethyst.commons.service.nwc.NwcPaymentTracker
import com.vitorpamplona.amethyst.commons.util.KmpLock
import com.vitorpamplona.amethyst.commons.util.withLock
import com.vitorpamplona.quartz.buzz.aeEngrams.EngramEvent
import com.vitorpamplona.quartz.buzz.agentProfiles.AgentProfileEvent
import com.vitorpamplona.quartz.buzz.amTurnMetrics.AgentTurnMetricEvent
@@ -424,6 +422,8 @@ import kotlinx.coroutines.flow.callbackFlow
import kotlinx.coroutines.flow.map
import kotlinx.coroutines.launch
import kotlin.concurrent.Volatile
import kotlin.concurrent.atomics.AtomicBoolean
import kotlin.concurrent.atomics.ExperimentalAtomicApi
import kotlin.time.TimeSource
/**
@@ -511,8 +511,8 @@ open class EventCache :
*/
val mintDirectory = MintDirectoryIndex()
@Volatile private var mintDirectoryBackfilled = false
private val mintDirectoryBackfillLock = KmpLock()
@OptIn(ExperimentalAtomicApi::class)
private val mintDirectoryBackfilled = AtomicBoolean(false)
/**
* Sweeps `notes` + `addressables` for any NIP-87 / NIP-61 event the
@@ -526,15 +526,17 @@ open class EventCache :
* scan failure is swallowed so the index stays usable even if the
* cache is in an unexpected state.
*/
@OptIn(ExperimentalAtomicApi::class)
fun ensureMintDirectoryBackfilled() {
if (mintDirectoryBackfilled) return
mintDirectoryBackfillLock.withLock {
if (mintDirectoryBackfilled) return
runCatching {
notes.forEach { _, note -> note.event?.let(::updateMintIndex) }
addressables.forEach { _, note -> note.event?.let(::updateMintIndex) }
}
mintDirectoryBackfilled = true
// Claim the sweep with a CAS rather than a lock: exactly one caller gets true and does
// the work. A second caller returns immediately instead of blocking for the length of a
// full cache scan — it may see a partially filled index, which is the same thing it sees
// before any backfill runs, and one relay round-trip later it does not.
if (!mintDirectoryBackfilled.compareAndSet(false, true)) return
runCatching {
notes.forEach { _, note -> note.event?.let(::updateMintIndex) }
addressables.forEach { _, note -> note.event?.let(::updateMintIndex) }
}
}
@@ -18,36 +18,160 @@
* 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.
*/
@file:OptIn(ExperimentalAtomicApi::class)
package com.vitorpamplona.amethyst.commons.model.observables
import com.vitorpamplona.amethyst.commons.model.AddressableNote
import com.vitorpamplona.amethyst.commons.model.Note
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 kotlin.concurrent.atomics.AtomicReference
import kotlin.concurrent.atomics.ExperimentalAtomicApi
/**
* A relay-like list of the events matching a filter, newest first, re-emitting on addressable updates.
* Creates a list of events (regular and addressable), sorted by created_at, that is updated
* every time a new matching event is received — INCLUDING newer versions of addressables, whose
* refreshed content re-emits (the entry keeps its captured position; consumers that care about
* order re-sort downstream).
*
* The declaration is shared so the event cache can live in `commonMain`; the implementation is
* not. One entry per idHex under concurrent ingest is kept by doing every write to the sorted
* set inside that key's `ConcurrentHashMap.compute` critical section — per-key striping, no
* lock — over a `ConcurrentSkipListSet` whose weakly-consistent iterator the snapshot
* de-duplicates. Kotlin/Native has neither primitive, and the argument for why this is correct
* (spelled out on the jvmAndroid actual) does not survive being reassembled out of weaker ones:
* a copy-on-write `compute` may re-run its lambda, and this lambda has side effects.
* There is exactly one [Note] instance per id/address (the cache owns their creation), so
* uniqueness is a non-issue in principle — except a note's sort key is mutable: a newer
* replaceable event swaps the event on the SAME [AddressableNote] instance, changing its
* created_at in place. Anything ordered on that live value cannot survive it. So the sort key is
* snapshotted into an immutable [Entry] when the note first enters and never read live again.
*
* So iOS has no implementation yet and every member throws — see `EventListMatchingFilter.ios.kt`.
* **Lock-free, by copy-on-write.** Observer callbacks fire concurrently from several consume
* threads (relay ingest + UI-side justConsume) and this observer is used everywhere, so it takes
* no lock: the whole list lives in one immutable [State] behind an [AtomicReference], and every
* mutation builds the next state and publishes it with a compare-and-set, retrying if another
* thread won the race. Same shape quartz uses for its native `ConcurrentMap`.
*
* Two things fall out of that, both previously paid for by hand:
* - a reader never sees a torn list, so the emitted snapshot needs no de-duplication pass —
* one entry per idHex holds by construction, since [State.ids] is what gates insertion;
* - the emitted list IS the published state, so emitting costs one map rather than a walk of a
* 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.
*/
expect class EventListMatchingFilter<T : Event>(
filter: Filter,
atOnce: (filter: Filter) -> List<Note>,
update: (List<T>) -> Unit,
class EventListMatchingFilter<T : Event>(
private val filter: Filter,
private val atOnce: (filter: Filter) -> List<Note>,
private val update: (List<T>) -> Unit,
) : Observable {
/** A note plus the sort key captured at insertion time, so ordering never depends on mutable state. */
private class Entry(
val note: Note,
val createdAt: Long,
val id: HexKey,
)
/** One immutable version of the list. [ids] is the membership gate, keyed by the stable idHex. */
private inner class State(
val entries: List<Entry>,
val ids: Set<HexKey>,
) {
@Suppress("UNCHECKED_CAST")
fun events(): List<T> = entries.mapNotNull { it.note.event as? T }
}
// created_at descending, id ascending as a stable tiebreak. Both fields are
// immutable snapshots, so an Entry never moves once inserted.
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 state = AtomicReference(State(emptyList(), emptySet()))
private fun entryFor(note: Note): Entry {
// A null event (unresolved note) sorts last, matching CreatedAtIdHexComparator.
val event = note.event
return Entry(note, note.createdAt() ?: Long.MIN_VALUE, event?.id ?: note.idHex)
}
private fun State.plus(
entry: Entry,
limit: Int?,
): State {
val found = entries.binarySearch(entry, order)
val at = if (found < 0) -found - 1 else found
val grown = ArrayList<Entry>(entries.size + 1)
grown.addAll(entries.subList(0, at))
grown.add(entry)
grown.addAll(entries.subList(at, entries.size))
if (limit == null || grown.size <= limit) return State(grown, ids + entry.note.idHex)
// 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)
}
private fun State.minus(idHex: HexKey): State = State(entries.filterNot { it.note.idHex == idHex }, ids - idHex)
override fun new(
event: Event,
note: Note,
)
) {
if (event is AddressableEvent && note !is AddressableNote) {
// The "version" note (a regular note holding an addressable event) is never
// stored — the AddressableNote is. Re-emit if that addressable is already
// listed, so consumers pick up the refreshed content read live off the note.
val current = state.load()
if (event.address().toValue() in current.ids) update(current.events())
return
}
override fun remove(note: Note)
if (!filter.match(event)) return
fun init()
val limit = filter.limit
while (true) {
val current = state.load()
// Already listed: the entry keeps its captured position — re-sorting it would
// mean a remove plus an add, and two entries for one note would transiently
// coexist and read the same live event twice. Re-emit so consumers see the
// refreshed content.
if (note.idHex in current.ids) {
update(current.events())
return
}
val next = current.plus(entryFor(note), limit)
if (state.compareAndSet(current, next)) {
update(next.events())
return
}
}
}
override fun remove(note: Note) {
while (true) {
val current = state.load()
if (note.idHex !in current.ids) return
val next = current.minus(note.idHex)
if (state.compareAndSet(current, next)) {
update(next.events())
return
}
}
}
fun init() {
// The cache query behind [atOnce] has already applied the filter's limit, so this
// inserts what it is given, exactly as the previous implementation did.
var fresh = State(emptyList(), emptySet())
atOnce(filter).forEach { note ->
if (note.idHex !in fresh.ids) fresh = fresh.plus(entryFor(note), null)
}
state.store(fresh)
update(fresh.events())
}
}
@@ -18,36 +18,151 @@
* 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.
*/
@file:OptIn(ExperimentalAtomicApi::class)
package com.vitorpamplona.amethyst.commons.model.observables
import com.vitorpamplona.amethyst.commons.model.AddressableNote
import com.vitorpamplona.amethyst.commons.model.Note
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 kotlin.concurrent.atomics.AtomicReference
import kotlin.concurrent.atomics.ExperimentalAtomicApi
/**
* A relay-like list of the notes matching a filter, newest first, that only grows.
* Creates a list of notes (regular and addressable), sorted by created_at like a
* relay, that only grows when a new note appears.
*
* The declaration is shared so the event cache can live in `commonMain`; the implementation is
* not. One entry per idHex under concurrent ingest is kept by doing every write to the sorted
* set inside that key's `ConcurrentHashMap.compute` critical section — per-key striping, no
* lock — over a `ConcurrentSkipListSet` whose weakly-consistent iterator the snapshot
* de-duplicates. Kotlin/Native has neither primitive, and the argument for why this is correct
* (spelled out on the jvmAndroid actual) does not survive being reassembled out of weaker ones:
* a copy-on-write `compute` may re-run its lambda, and this lambda has side effects.
* New versions of addressables do not update the list.
*
* So iOS has no implementation yet and every member throws — see `NoteListMatchingFilter.ios.kt`.
* There is exactly one [Note] instance per id/address (the cache owns their creation), so
* uniqueness is a non-issue in principle — except a note's sort key is mutable: a newer
* replaceable event swaps the event on the SAME [AddressableNote] instance, changing its
* created_at in place. Anything ordered on that live value cannot survive it. So the sort key is
* snapshotted into an immutable [Entry] when the note first enters and never read live again.
*
* **Lock-free, by copy-on-write.** Observer callbacks fire concurrently from several consume
* threads (relay ingest + UI-side justConsume) and this observer is used everywhere, so it takes
* no lock: the whole list lives in one immutable [State] behind an [AtomicReference], and every
* mutation builds the next state and publishes it with a compare-and-set, retrying if another
* thread won the race. Same shape quartz uses for its native `ConcurrentMap`.
*
* Two things fall out of that, both previously paid for by hand:
* - a reader never sees a torn list, so the emitted snapshot needs no de-duplication pass —
* one entry per idHex holds by construction, since [State.ids] is what gates insertion;
* - the emitted list IS the published state, so emitting costs one map rather than a walk of a
* 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.
*/
expect class NoteListMatchingFilter(
filter: Filter,
atOnce: (filter: Filter) -> List<Note>,
update: (List<Note>) -> Unit,
class NoteListMatchingFilter(
private val filter: Filter,
private val atOnce: (filter: Filter) -> List<Note>,
private val update: (List<Note>) -> Unit,
) : Observable {
/** A note plus the sort key captured at insertion time, so ordering never depends on mutable state. */
private class Entry(
val note: Note,
val createdAt: Long,
val id: HexKey,
)
/** One immutable version of the list. [ids] is the membership gate, keyed by the stable idHex. */
private class State(
val entries: List<Entry>,
val ids: Set<HexKey>,
) {
fun notes(): List<Note> = entries.map { it.note }
companion object {
val EMPTY = State(emptyList(), emptySet())
}
}
// created_at descending, id ascending as a stable tiebreak. Both fields are
// immutable snapshots, so an Entry never moves once inserted.
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 state = AtomicReference(State.EMPTY)
private fun entryFor(note: Note): Entry {
// A null event (unresolved note) sorts last, matching CreatedAtIdHexComparator.
val event = note.event
return Entry(note, note.createdAt() ?: Long.MIN_VALUE, event?.id ?: note.idHex)
}
private fun State.plus(
entry: Entry,
limit: Int?,
): State {
val found = entries.binarySearch(entry, order)
val at = if (found < 0) -found - 1 else found
val grown = ArrayList<Entry>(entries.size + 1)
grown.addAll(entries.subList(0, at))
grown.add(entry)
grown.addAll(entries.subList(at, entries.size))
if (limit == null || grown.size <= limit) return State(grown, ids + entry.note.idHex)
// 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)
}
private fun State.minus(idHex: HexKey): State = State(entries.filterNot { it.note.idHex == idHex }, ids - idHex)
override fun new(
event: Event,
note: Note,
)
) {
if (event is AddressableEvent && note !is AddressableNote) return
override fun remove(note: Note)
if (!filter.match(event)) return
fun init()
val limit = filter.limit
while (true) {
val current = state.load()
// Already listed: a newer version of an addressable keeps its entry, and its
// position with it, so the list only ever grows.
if (note.idHex in current.ids) return
val next = current.plus(entryFor(note), limit)
if (state.compareAndSet(current, next)) {
update(next.notes())
return
}
}
}
override fun remove(note: Note) {
while (true) {
val current = state.load()
if (note.idHex !in current.ids) return
val next = current.minus(note.idHex)
if (state.compareAndSet(current, next)) {
update(next.notes())
return
}
}
}
fun init() {
// The cache query behind [atOnce] has already applied the filter's limit, so this
// inserts what it is given, exactly as the previous implementation did.
var fresh = State.EMPTY
atOnce(filter).forEach { note ->
if (note.idHex !in fresh.ids) fresh = fresh.plus(entryFor(note), null)
}
state.store(fresh)
update(fresh.notes())
}
}
@@ -1,52 +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.amethyst.commons.model.observables
import com.vitorpamplona.amethyst.commons.model.Note
import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
actual class EventListMatchingFilter<T : Event> actual constructor(
filter: Filter,
atOnce: (filter: Filter) -> List<Note>,
update: (List<T>) -> Unit,
) : Observable {
actual override fun new(
event: Event,
note: Note,
) {
notYetImplemented()
}
actual override fun remove(note: Note) {
notYetImplemented()
}
actual fun init() {
notYetImplemented()
}
}
private fun notYetImplemented(): Nothing =
throw NotImplementedError(
"The matching-filter observables have no iOS implementation yet: they need a per-key " +
"atomic section over a sorted concurrent set.",
)
@@ -1,52 +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.amethyst.commons.model.observables
import com.vitorpamplona.amethyst.commons.model.Note
import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
actual class NoteListMatchingFilter actual constructor(
filter: Filter,
atOnce: (filter: Filter) -> List<Note>,
update: (List<Note>) -> Unit,
) : Observable {
actual override fun new(
event: Event,
note: Note,
) {
notYetImplemented()
}
actual override fun remove(note: Note) {
notYetImplemented()
}
actual fun init() {
notYetImplemented()
}
}
private fun notYetImplemented(): Nothing =
throw NotImplementedError(
"The matching-filter observables have no iOS implementation yet: they need a per-key " +
"atomic section over a sorted concurrent set.",
)
@@ -1,167 +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.amethyst.commons.model.observables
import com.vitorpamplona.amethyst.commons.model.AddressableNote
import com.vitorpamplona.amethyst.commons.model.Note
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 java.util.concurrent.ConcurrentHashMap
import java.util.concurrent.ConcurrentSkipListSet
/**
* Creates a list of events (regular and addressable), sorted by created_at, that
* is updated every time a new matching event is received — INCLUDING newer
* versions of addressables, whose refreshed content re-emits and re-sorts.
*
* Like [NoteListMatchingFilter], this cannot store mutable [Note]s in a set
* ordered on their live created_at: a newer replaceable event mutates the SAME
* [AddressableNote] instance in place, which strands its node and lets the same
* instance be inserted twice — the emitted list would then carry the same event
* twice. So the sort key is snapshotted into an immutable [Entry], [sorted] is
* ordered on that snapshot, and [byId] is the membership source of truth keyed
* by the stable idHex (the address for addressables, the event id otherwise).
*
* Unlike [NoteListMatchingFilter], an addressable update is NOT ignored: the
* list re-emits so consumers pick up the refreshed [Event] (read live off the
* note). The entry's captured position is kept — re-sorting an updated entry
* would mean a remove+add on [sorted], and two entries with different captured
* keys for the same note would then transiently coexist and both read the same
* live event, duplicating it. Keeping the position matches the original's only
* non-corrupting behavior (it never reliably re-sorted either); consumers that
* care about order re-sort downstream.
*
* Observer callbacks fire concurrently from several consume threads (relay
* ingest + UI-side justConsume), so it stays lock-free: [sorted] is only ever
* written for a key INSIDE that key's `compute` critical section, and only while
* the key is absent. ConcurrentHashMap stripes per key, so same-idHex ops
* serialize while different keys run in parallel. An entry is added only while
* its key is absent from [byId], and every path that frees a key removes its
* entry from [sorted] first.
*
* That keeps [sorted] converged to one entry per idHex, but a
* ConcurrentSkipListSet iterator is only weakly consistent: under concurrent
* add/remove churn a single traversal can momentarily surface a key twice
* (lazy-deleted node not yet unlinked while its replacement is inserted). The
* emitted list must never carry a duplicate id — a LazyColumn keyed on it would
* crash — so [snapshot] deduplicates by the stable idHex as it materializes.
*/
actual class EventListMatchingFilter<T : Event> actual constructor(
private val filter: Filter,
private val atOnce: (filter: Filter) -> List<Note>,
private val update: (List<T>) -> Unit,
) : Observable {
/** A note plus the sort key captured at insertion time, so ordering never depends on mutable state. */
private class Entry(
val note: Note,
val createdAt: Long,
val id: HexKey,
)
// created_at descending, id ascending as a stable tiebreak. Both fields are
// immutable snapshots, so an Entry never moves once inserted.
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 {
val event = note.event
return Entry(note, note.createdAt() ?: Long.MIN_VALUE, event?.id ?: note.idHex)
}
@Suppress("UNCHECKED_CAST")
actual override fun new(
event: Event,
note: Note,
) {
if (event is AddressableEvent && note !is AddressableNote) {
// The "version" note (a regular note holding an addressable event) is
// never stored — the AddressableNote is. Re-emit if that addressable
// is already listed so consumers pick up the refreshed content.
if (byId.containsKey(event.address().toValue())) update(snapshot())
return
}
if (!filter.match(event)) return
// Add to [sorted] atomically with claiming the idHex slot, only when the
// key is absent. An update keeps its entry (and position) — the re-emit
// below reflects the refreshed event read live off the note.
var added = false
byId.compute(note.idHex) { _, existing ->
existing ?: entryFor(note).also {
sorted.add(it)
added = true
}
}
if (added) {
val limit = filter.limit
if (limit != null && sorted.size > limit) {
// Drop the oldest (sorts last under [order]).
sorted.pollLast()?.let { byId.remove(it.note.idHex, it) }
}
}
// Always re-emit on a match: a first insert grows the list, an update
// refreshes the event content the snapshot reads off the note.
update(snapshot())
}
@Suppress("UNCHECKED_CAST")
actual 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())
}
@Suppress("UNCHECKED_CAST")
actual fun init() {
sorted.clear()
byId.clear()
atOnce(filter).forEach { note ->
byId.computeIfAbsent(note.idHex) { entryFor(note).also { sorted.add(it) } }
}
update(snapshot())
}
@Suppress("UNCHECKED_CAST")
private fun snapshot(): List<T> {
// Dedup by the stable idHex: the weakly-consistent iterator can transiently
// surface a key twice under concurrent churn. Both would read the same live
// event off the same note, so keeping the first (newest position) is correct.
val seen = HashSet<HexKey>()
return sorted.mapNotNull { e -> if (seen.add(e.note.idHex)) e.note.event as? T else null }
}
}
@@ -1,154 +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.amethyst.commons.model.observables
import com.vitorpamplona.amethyst.commons.model.AddressableNote
import com.vitorpamplona.amethyst.commons.model.Note
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 java.util.concurrent.ConcurrentHashMap
import java.util.concurrent.ConcurrentSkipListSet
/**
* Creates a list of notes (regular and addressable), sorted by created_at like a
* relay, that only grows when a new note appears.
*
* New versions of addressables do not update the list.
*
* There is exactly one [Note] instance per id/address (LocalCache owns their
* creation), so uniqueness is a non-issue in principle — except a note's sort
* key is mutable: a newer replaceable event swaps the event on the SAME
* [AddressableNote] instance, changing its created_at in place. A sorted set
* ordered on that live value cannot survive it — the moved node is no longer on
* the search path of add()/remove(), so the same instance gets inserted twice
* and the emitted list carries a duplicate idHex, crashing any LazyColumn keyed
* on it.
*
* So the sort key is snapshotted into an immutable [Entry] when the note first
* enters and never read live again; [sorted] is ordered on that snapshot
* (stable) and [byId] is the membership source of truth, keyed by the stable
* idHex.
*
* Observer callbacks fire concurrently from several consume threads (relay
* ingest + UI-side justConsume), and this observer is used everywhere, so it
* stays lock-free: [byId] is a ConcurrentHashMap and every write to [sorted] for
* a given idHex happens INSIDE that key's `compute` critical section.
* ConcurrentHashMap stripes per key, so same-idHex ops serialize while different
* keys run fully in parallel. An entry is added to [sorted] only while its key
* is absent from [byId], and every path that makes a key absent removes its
* entry from [sorted] first, keeping [sorted] converged to one entry per idHex.
*
* That convergence isn't enough on its own: a ConcurrentSkipListSet iterator is
* only weakly consistent, so under concurrent add/remove churn a single
* traversal can momentarily surface a key twice (a lazy-deleted node not yet
* unlinked while its replacement is inserted). The emitted list must never carry
* a duplicate idHex — the LazyColumn keyed on it would crash — so [snapshot]
* deduplicates by idHex as it materializes.
*/
actual class NoteListMatchingFilter actual constructor(
private val filter: Filter,
private val atOnce: (filter: Filter) -> List<Note>,
private val update: (List<Note>) -> Unit,
) : Observable {
/** A note plus the sort key captured at insertion time, so ordering never depends on mutable state. */
private class Entry(
val note: Note,
val createdAt: Long,
val id: HexKey,
)
// created_at descending, id ascending as a stable tiebreak. Both fields are
// immutable snapshots, so an Entry never moves once inserted.
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 {
// A null event (unresolved note) sorts last, matching CreatedAtIdHexComparator.
val event = note.event
return Entry(note, note.createdAt() ?: Long.MIN_VALUE, event?.id ?: note.idHex)
}
actual override fun new(
event: Event,
note: Note,
) {
if (event is AddressableEvent && note !is AddressableNote) return
if (!filter.match(event)) return
// Add to [sorted] atomically with claiming the idHex slot. New versions
// of an already listed note return the existing entry untouched.
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) {
// Drop the oldest (sorts last under [order]).
sorted.pollLast()?.let { byId.remove(it.note.idHex, it) }
}
update(snapshot())
}
actual override fun remove(note: Note) {
// Remove from [sorted] atomically with releasing the idHex slot.
var removed = false
byId.compute(note.idHex) { _, existing ->
if (existing != null) {
sorted.remove(existing)
removed = true
}
null
}
if (removed) update(snapshot())
}
actual fun init() {
sorted.clear()
byId.clear()
atOnce(filter).forEach { note ->
byId.computeIfAbsent(note.idHex) { entryFor(note).also { sorted.add(it) } }
}
update(snapshot())
}
private fun snapshot(): List<Note> {
// Dedup by idHex: the weakly-consistent iterator can transiently surface a
// key twice under concurrent churn. Keeping the first (newest position) is
// correct — both nodes point at the same note.
val seen = HashSet<HexKey>()
return sorted.mapNotNull { e -> e.note.takeIf { seen.add(it.idHex) } }
}
}