From c25cf6524f20a1a15cbf7f5d8c90536fbe32e1dc Mon Sep 17 00:00:00 2001 From: Claude Date: Mon, 21 Sep 2026 17:57:14 +0000 Subject: [PATCH] refactor(commons): make the list observables lock-free and common MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 Claude-Session: https://claude.ai/code/session_01KkULS5SVq4GHDdoCzajKi8 --- .../commons/model/cache/EventCache.kt | 26 +-- .../observables/EventListMatchingFilter.kt | 156 ++++++++++++++-- .../observables/NoteListMatchingFilter.kt | 147 +++++++++++++-- .../EventListMatchingFilter.ios.kt | 52 ------ .../observables/NoteListMatchingFilter.ios.kt | 52 ------ .../EventListMatchingFilter.jvmAndroid.kt | 167 ------------------ .../NoteListMatchingFilter.jvmAndroid.kt | 154 ---------------- 7 files changed, 285 insertions(+), 469 deletions(-) delete mode 100644 commons/src/iosMain/kotlin/com/vitorpamplona/amethyst/commons/model/observables/EventListMatchingFilter.ios.kt delete mode 100644 commons/src/iosMain/kotlin/com/vitorpamplona/amethyst/commons/model/observables/NoteListMatchingFilter.ios.kt delete mode 100644 commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/model/observables/EventListMatchingFilter.jvmAndroid.kt delete mode 100644 commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/model/observables/NoteListMatchingFilter.jvmAndroid.kt diff --git a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/cache/EventCache.kt b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/cache/EventCache.kt index 5ffd5bb72f..85ca2ebe0e 100644 --- a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/cache/EventCache.kt +++ b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/cache/EventCache.kt @@ -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) } } } diff --git a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/observables/EventListMatchingFilter.kt b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/observables/EventListMatchingFilter.kt index 921cad683e..dba6c0eb20 100644 --- a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/observables/EventListMatchingFilter.kt +++ b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/observables/EventListMatchingFilter.kt @@ -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( - filter: Filter, - atOnce: (filter: Filter) -> List, - update: (List) -> Unit, +class EventListMatchingFilter( + private val filter: Filter, + private val atOnce: (filter: Filter) -> List, + private val update: (List) -> 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, + val ids: Set, + ) { + @Suppress("UNCHECKED_CAST") + fun events(): List = 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 { 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(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()) + } } diff --git a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/observables/NoteListMatchingFilter.kt b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/observables/NoteListMatchingFilter.kt index f3811c54ec..bd8603729b 100644 --- a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/observables/NoteListMatchingFilter.kt +++ b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/observables/NoteListMatchingFilter.kt @@ -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, - update: (List) -> Unit, +class NoteListMatchingFilter( + private val filter: Filter, + private val atOnce: (filter: Filter) -> List, + private val update: (List) -> 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, + val ids: Set, + ) { + fun notes(): List = 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 { 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(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()) + } } diff --git a/commons/src/iosMain/kotlin/com/vitorpamplona/amethyst/commons/model/observables/EventListMatchingFilter.ios.kt b/commons/src/iosMain/kotlin/com/vitorpamplona/amethyst/commons/model/observables/EventListMatchingFilter.ios.kt deleted file mode 100644 index 37831da78e..0000000000 --- a/commons/src/iosMain/kotlin/com/vitorpamplona/amethyst/commons/model/observables/EventListMatchingFilter.ios.kt +++ /dev/null @@ -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 actual constructor( - filter: Filter, - atOnce: (filter: Filter) -> List, - update: (List) -> 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.", - ) diff --git a/commons/src/iosMain/kotlin/com/vitorpamplona/amethyst/commons/model/observables/NoteListMatchingFilter.ios.kt b/commons/src/iosMain/kotlin/com/vitorpamplona/amethyst/commons/model/observables/NoteListMatchingFilter.ios.kt deleted file mode 100644 index be2bb3e51f..0000000000 --- a/commons/src/iosMain/kotlin/com/vitorpamplona/amethyst/commons/model/observables/NoteListMatchingFilter.ios.kt +++ /dev/null @@ -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, - update: (List) -> 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.", - ) diff --git a/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/model/observables/EventListMatchingFilter.jvmAndroid.kt b/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/model/observables/EventListMatchingFilter.jvmAndroid.kt deleted file mode 100644 index 0e5a6ec4cc..0000000000 --- a/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/model/observables/EventListMatchingFilter.jvmAndroid.kt +++ /dev/null @@ -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 actual constructor( - private val filter: Filter, - private val atOnce: (filter: Filter) -> List, - private val update: (List) -> 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 { 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() - - 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 { - // 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() - return sorted.mapNotNull { e -> if (seen.add(e.note.idHex)) e.note.event as? T else null } - } -} diff --git a/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/model/observables/NoteListMatchingFilter.jvmAndroid.kt b/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/model/observables/NoteListMatchingFilter.jvmAndroid.kt deleted file mode 100644 index f3466a548d..0000000000 --- a/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/model/observables/NoteListMatchingFilter.jvmAndroid.kt +++ /dev/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, - private val update: (List) -> 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 { 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() - - 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 { - // 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() - return sorted.mapNotNull { e -> e.note.takeIf { seen.add(it.idHex) } } - } -}