diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/DebugUtils.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/DebugUtils.kt index 6b2880415..5869a936a 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/DebugUtils.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/DebugUtils.kt @@ -139,7 +139,7 @@ fun debugState(context: Context) { Log.d( STATE_DUMP_TAG, "Observables: " + - LocalCache.observables.size, + LocalCache.observables.size(), ) Log.d( diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/model/LocalCache.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/model/LocalCache.kt index f6bfbb53a..9bd2c6198 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/model/LocalCache.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/model/LocalCache.kt @@ -88,6 +88,7 @@ import com.vitorpamplona.quartz.nip01Core.metadata.MetadataEvent import com.vitorpamplona.quartz.nip01Core.relay.client.single.IRelayClient import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.EventCmd import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip01Core.relay.filters.FilterIndex import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl import com.vitorpamplona.quartz.nip01Core.tags.aTag.ATag import com.vitorpamplona.quartz.nip01Core.tags.aTag.taggedAddresses @@ -261,7 +262,6 @@ import java.io.File import java.io.FileOutputStream import java.io.IOException import java.util.SortedSet -import java.util.concurrent.ConcurrentHashMap interface ILocalCache { fun markAsSeen( @@ -290,7 +290,15 @@ object LocalCache : ILocalCache, ICacheProvider { val deletionIndex = DeletionIndex() - val observables = ConcurrentHashMap(10) + /** + * Inverted index over the active [Observable]s. New events fan + * out only to observers whose filter actually narrows on a field + * the event carries (author, kind, single-letter tag), instead + * of waking every observer per event. Predicate-only observers + * (no underlying [Filter]) live in the unindexed pool and still + * see every event. + */ + val observables = FilterIndex() fun Filter.match(note: Note): Boolean { val event = note.event @@ -361,10 +369,10 @@ object LocalCache : ILocalCache, ICacheProvider { newFilter.init() - observables[newFilter] = newFilter + observables.register(filter, newFilter) awaitClose { - observables.remove(newFilter) + observables.unregister(newFilter) } }.buffer(kotlinx.coroutines.channels.Channel.CONFLATED) @@ -377,10 +385,10 @@ object LocalCache : ILocalCache, ICacheProvider { cachedFilter.init() - observables.put(cachedFilter, cachedFilter) + observables.register(filter, cachedFilter) awaitClose { - observables.remove(cachedFilter) + observables.unregister(cachedFilter) } }.buffer(kotlinx.coroutines.channels.Channel.CONFLATED) @@ -402,14 +410,28 @@ object LocalCache : ILocalCache, ICacheProvider { trySend(it) } - observables.put(newFilter, newFilter) + // Unindexed: predicate is opaque, the index can't narrow. + // Caller is delivered every event and runs the predicate. + observables.registerUnindexed(newFilter) awaitClose { - observables.remove(newFilter) + observables.unregister(newFilter) } } - fun observeNewEvents(filter: Filter): Flow = observeNewEvents(filter::match) + fun observeNewEvents(filter: Filter): Flow = + callbackFlow { + val newFilter = + NewEventMatchingFilter(filter::match) { + trySend(it) + } + + observables.register(filter, newFilter) + + awaitClose { + observables.unregister(newFilter) + } + } @Suppress("UNCHECKED_CAST") fun observeLatestEvent(filter: Filter) = observeEvents(filter).map { it.firstOrNull() } @@ -2514,22 +2536,21 @@ object LocalCache : ILocalCache, ICacheProvider { private fun refreshNewNoteObservers(newNote: Note) { val event = newNote.event as Event - val observableBiConsumer = - java.util.function.BiConsumer { _, u -> - u.new(event, newNote) - } - - observables.forEach(observableBiConsumer) + // Index-driven fanout: only observers whose filter narrows + // on a field this event carries (or that registered as + // unindexed) get woken up. The match check inside each + // observer's `new()` still enforces negative constraints. + for (observer in observables.candidatesFor(event)) { + observer.new(event, newNote) + } live.newNote(newNote) } private fun refreshDeletedNoteObservers(newNote: Note) { - val observableBiConsumer = - java.util.function.BiConsumer { _, u -> - u.remove(newNote) - } - - observables.forEach(observableBiConsumer) + // Deletes don't have a filterable shape — every observer + // might hold this note in its result set, so iterate them + // all. The index doesn't help here. + observables.forEach { it.remove(newNote) } live.removedNote(newNote) }