Renaming and making a test case that considers limits.

Fixes active filter index bug
This commit is contained in:
Vitor Pamplona
2026-03-17 15:45:52 -04:00
parent b8351b8e29
commit 3e19195e4a
6 changed files with 277 additions and 133 deletions
@@ -23,7 +23,7 @@ package com.vitorpamplona.amethyst.ui.screen.loggedIn.relays.eventsync
import com.vitorpamplona.amethyst.model.Account
import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.core.HexKey
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.downloadFromRelay
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.reqBypassingRelayLimits
import com.vitorpamplona.quartz.nip01Core.relay.client.listeners.IRelayClientListener
import com.vitorpamplona.quartz.nip01Core.relay.client.single.IRelayClient
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.Message
@@ -126,9 +126,9 @@ class EventSync(
*/
data class LiveSyncActivity(
val recentCompletions: List<CompletedRelayInfo> = emptyList(),
val outboxTargets: Set<NormalizedRelayUrl> = emptySet(),
val inboxTargets: Set<NormalizedRelayUrl> = emptySet(),
val dmTargets: Set<NormalizedRelayUrl> = emptySet(),
val outboxTargets: List<DestinationRelayInfo> = emptyList(),
val inboxTargets: List<DestinationRelayInfo> = emptyList(),
val dmTargets: List<DestinationRelayInfo> = emptyList(),
) {
/**
* @param eventsFound Total events received from this relay across all pages.
@@ -141,6 +141,17 @@ class EventSync(
val eventsFound: Int,
val eventsAccepted: Int,
)
/**
* @param relay The destination relay URL.
* @param eventsSent Number of events sent to this relay.
* @param eventsAccepted Number of OK=true responses received from this relay.
*/
data class DestinationRelayInfo(
val relay: NormalizedRelayUrl,
val eventsSent: Int,
val eventsAccepted: Int,
)
}
private val _syncState = MutableStateFlow<SyncState>(SyncState.Idle)
@@ -163,14 +174,28 @@ class EventSync(
@Volatile private var liveDmTargets: Set<NormalizedRelayUrl> = emptySet()
/** Per-destination-relay counters, updated atomically during a sync run. */
private val liveSentCountPerDestRelay = ConcurrentHashMap<NormalizedRelayUrl, AtomicInteger>()
private val liveAcceptedCountPerDestRelay = ConcurrentHashMap<NormalizedRelayUrl, AtomicInteger>()
private fun emitLiveSnapshot() {
val completions = synchronized(trackingLock) { trackingCompletions.toList() }
fun buildDestInfo(relays: Set<NormalizedRelayUrl>) =
relays.map { relay ->
LiveSyncActivity.DestinationRelayInfo(
relay = relay,
eventsSent = liveSentCountPerDestRelay[relay]?.get() ?: 0,
eventsAccepted = liveAcceptedCountPerDestRelay[relay]?.get() ?: 0,
)
}
_liveActivity.value =
LiveSyncActivity(
recentCompletions = completions,
outboxTargets = liveOutboxTargets,
inboxTargets = liveInboxTargets,
dmTargets = liveDmTargets,
outboxTargets = buildDestInfo(liveOutboxTargets),
inboxTargets = buildDestInfo(liveInboxTargets),
dmTargets = buildDestInfo(liveDmTargets),
)
}
@@ -183,6 +208,8 @@ class EventSync(
fun start() {
if (_syncState.value is SyncState.Running) return
synchronized(trackingLock) { trackingCompletions.clear() }
liveSentCountPerDestRelay.clear()
liveAcceptedCountPerDestRelay.clear()
_liveActivity.value = LiveSyncActivity()
syncJob =
scope.launch(Dispatchers.IO) {
@@ -294,6 +321,7 @@ class EventSync(
val sourceRelay = sourceRelayOfEvent.remove(msg.eventId) ?: return
acceptedCountPerRelay.getOrPut(sourceRelay) { AtomicInteger(0) }.incrementAndGet()
totalAccepted.incrementAndGet()
liveAcceptedCountPerDestRelay.getOrPut(relay.url) { AtomicInteger(0) }.incrementAndGet()
}
}
}
@@ -318,6 +346,9 @@ class EventSync(
sourceRelayOfEvent[event.id] = sourceRelay
account.client.send(event, outboxTargets)
totalSent.incrementAndGet()
outboxTargets.forEach { dest ->
liveSentCountPerDestRelay.getOrPut(dest) { AtomicInteger(0) }.incrementAndGet()
}
}
}
val pTagsMe = event.tags.isTaggedUser(myPubKey)
@@ -327,12 +358,18 @@ class EventSync(
sourceRelayOfEvent[event.id] = sourceRelay
account.client.send(event, dmTargets)
totalSent.incrementAndGet()
dmTargets.forEach { dest ->
liveSentCountPerDestRelay.getOrPut(dest) { AtomicInteger(0) }.incrementAndGet()
}
}
} else {
if (inboxTargets.isNotEmpty() && inboxSent.add(event.id)) {
sourceRelayOfEvent[event.id] = sourceRelay
account.client.send(event, inboxTargets)
totalSent.incrementAndGet()
inboxTargets.forEach { dest ->
liveSentCountPerDestRelay.getOrPut(dest) { AtomicInteger(0) }.incrementAndGet()
}
}
}
}
@@ -410,5 +447,5 @@ class EventSync(
relay: NormalizedRelayUrl,
baseFilters: List<Filter>,
onEvent: (Event) -> Unit,
): Int = account.client.downloadFromRelay(relay, baseFilters, RELAY_TIMEOUT_MS, onEvent)
): Int = account.client.reqBypassingRelayLimits(relay, baseFilters, RELAY_TIMEOUT_MS, onEvent)
}
@@ -24,8 +24,6 @@ import androidx.compose.foundation.background
import androidx.compose.foundation.layout.Arrangement
import androidx.compose.foundation.layout.Box
import androidx.compose.foundation.layout.Column
import androidx.compose.foundation.layout.ExperimentalLayoutApi
import androidx.compose.foundation.layout.FlowRow
import androidx.compose.foundation.layout.Row
import androidx.compose.foundation.layout.Spacer
import androidx.compose.foundation.layout.fillMaxSize
@@ -47,8 +45,6 @@ import androidx.compose.material3.LinearProgressIndicator
import androidx.compose.material3.MaterialTheme
import androidx.compose.material3.OutlinedButton
import androidx.compose.material3.Scaffold
import androidx.compose.material3.SuggestionChip
import androidx.compose.material3.SuggestionChipDefaults
import androidx.compose.material3.Text
import androidx.compose.material3.TextButton
import androidx.compose.runtime.Composable
@@ -452,11 +448,10 @@ private fun DestinationRelaysCard(activity: EventSync.LiveSyncActivity) {
}
}
@OptIn(ExperimentalLayoutApi::class)
@Composable
private fun DestinationSection(
label: String,
relays: Set<NormalizedRelayUrl>,
relays: List<EventSync.LiveSyncActivity.DestinationRelayInfo>,
color: androidx.compose.ui.graphics.Color,
) {
Text(
@@ -466,26 +461,58 @@ private fun DestinationSection(
color = MaterialTheme.colorScheme.onSurfaceVariant,
)
Spacer(Modifier.height(6.dp))
FlowRow(
horizontalArrangement = Arrangement.spacedBy(8.dp),
verticalArrangement = Arrangement.spacedBy(6.dp),
Column(verticalArrangement = Arrangement.spacedBy(2.dp)) {
relays.forEachIndexed { index, info ->
if (index > 0) {
HorizontalDivider(
color = MaterialTheme.colorScheme.outlineVariant.copy(alpha = 0.5f),
)
}
DestinationRelayRow(info = info, color = color)
}
}
}
@Composable
private fun DestinationRelayRow(
info: EventSync.LiveSyncActivity.DestinationRelayInfo,
color: androidx.compose.ui.graphics.Color,
) {
Row(
modifier =
Modifier
.fillMaxWidth()
.padding(vertical = 6.dp),
verticalAlignment = Alignment.CenterVertically,
horizontalArrangement = Arrangement.spacedBy(10.dp),
) {
relays.forEach { relay ->
SuggestionChip(
onClick = {},
label = {
Text(
text = relay.displayHost(),
style = MaterialTheme.typography.labelSmall,
maxLines = 1,
overflow = TextOverflow.Ellipsis,
)
},
colors =
SuggestionChipDefaults.suggestionChipColors(
containerColor = color.copy(alpha = 0.12f),
labelColor = color,
),
Box(
modifier =
Modifier
.size(8.dp)
.clip(CircleShape)
.background(color),
)
Text(
text = info.relay.displayHost(),
style = MaterialTheme.typography.bodySmall,
color = MaterialTheme.colorScheme.onSurface,
maxLines = 1,
overflow = TextOverflow.Ellipsis,
modifier = Modifier.weight(1f),
)
if (info.eventsSent > 0) {
Text(
text = stringRes(R.string.event_sync_log_recv, formatCount(info.eventsSent)),
style = MaterialTheme.typography.bodySmall,
color = MaterialTheme.colorScheme.onSurfaceVariant,
)
Spacer(Modifier.width(8.dp))
Text(
text = stringRes(R.string.event_sync_log_new, formatCount(info.eventsAccepted)),
style = MaterialTheme.typography.bodySmall,
fontWeight = if (info.eventsAccepted > 0) FontWeight.SemiBold else FontWeight.Normal,
color = if (info.eventsAccepted > 0) color else MaterialTheme.colorScheme.onSurfaceVariant,
)
}
}
@@ -658,18 +685,18 @@ private val previewActivity =
EventSync.LiveSyncActivity(
recentCompletions = previewCompletions,
outboxTargets =
setOf(
NormalizedRelayUrl("wss://outbox.nostr.com"),
NormalizedRelayUrl("wss://relay.damus.io"),
listOf(
EventSync.LiveSyncActivity.DestinationRelayInfo(NormalizedRelayUrl("wss://outbox.nostr.com"), 1247, 891),
EventSync.LiveSyncActivity.DestinationRelayInfo(NormalizedRelayUrl("wss://relay.damus.io"), 892, 45),
),
inboxTargets =
setOf(
NormalizedRelayUrl("wss://inbox.nostr.com"),
NormalizedRelayUrl("wss://nos.lol"),
listOf(
EventSync.LiveSyncActivity.DestinationRelayInfo(NormalizedRelayUrl("wss://inbox.nostr.com"), 500, 500),
EventSync.LiveSyncActivity.DestinationRelayInfo(NormalizedRelayUrl("wss://nos.lol"), 0, 0),
),
dmTargets =
setOf(
NormalizedRelayUrl("wss://dm.nostr.com"),
listOf(
EventSync.LiveSyncActivity.DestinationRelayInfo(NormalizedRelayUrl("wss://dm.nostr.com"), 15, 10),
),
)
@@ -34,7 +34,7 @@ import kotlinx.coroutines.withTimeoutOrNull
import kotlin.coroutines.coroutineContext
/**
* Downloads all pages of events matching [baseFilters] from a single [relay] using
* Downloads all pages of events matching [filters] from a single [relay] using
* paginated `until` cursors.
*
* After EOSE the oldest [Event.createdAt] seen in that page minus one becomes the
@@ -46,14 +46,14 @@ import kotlin.coroutines.coroutineContext
* Filters without a limit are considered unbounded and only stop on empty pages.
*
* @param relay The relay to query.
* @param baseFilters Filters to apply on every page (the `until` field is overwritten per page).
* @param filters Filters to apply on every page (the `until` field is overwritten per page).
* @param timeoutMs Maximum time to wait for a single page's EOSE before giving up.
* @param onEvent Called for every event received (in page order, after each EOSE).
* @return Total number of events received across all pages.
*/
suspend fun INostrClient.downloadFromRelay(
suspend fun INostrClient.reqBypassingRelayLimits(
relay: NormalizedRelayUrl,
baseFilters: List<Filter>,
filters: List<Filter>,
timeoutMs: Long = 30_000L,
onEvent: (Event) -> Unit,
): Int {
@@ -61,31 +61,29 @@ suspend fun INostrClient.downloadFromRelay(
var totalEvents = 0
// Track how many matching events each filter has received so far.
val matchCountPerFilter = IntArray(baseFilters.size)
val matchCountPerFilter = IntArray(filters.size)
while (true) {
coroutineContext.ensureActive()
// Only include filters that still need more events.
val activeFilterIndices =
baseFilters.indices.filter { i ->
val limit = baseFilters[i].limit
limit == null || matchCountPerFilter[i] < limit
val remainingFilters =
filters.filterIndexed { index, filter ->
val limit = filter.limit
limit == null || matchCountPerFilter[index] < limit
}
if (activeFilterIndices.isEmpty()) break
val activeBaseFilters = activeFilterIndices.map { baseFilters[it] }
if (remainingFilters.isEmpty()) break
val eventChannel = Channel<Event>(UNLIMITED)
val doneChannel = Channel<Unit>(Channel.CONFLATED)
val subId = newSubId()
val filters =
val activeFilters =
if (until == null) {
activeBaseFilters
remainingFilters
} else {
activeBaseFilters.map { it.copy(until = until) }
remainingFilters.map { it.copy(until = until) }
}
val listener =
@@ -119,11 +117,12 @@ suspend fun INostrClient.downloadFromRelay(
message: String,
forFilters: List<Filter>?,
) {
println("AABBCC $message")
doneChannel.trySend(Unit)
}
}
openReqSubscription(subId, mapOf(relay to filters), listener)
openReqSubscription(subId, mapOf(relay to activeFilters), listener)
withTimeoutOrNull(timeoutMs) { doneChannel.receive() }
close(subId)
eventChannel.close()
@@ -137,9 +136,15 @@ suspend fun INostrClient.downloadFromRelay(
if (event.createdAt < pageMinTs) pageMinTs = event.createdAt
// Count this event against every base filter it matches.
for (i in baseFilters.indices) {
if (baseFilters[i].match(event)) {
matchCountPerFilter[i]++
if (matchCountPerFilter.size == 1) {
// no need to run the match.
matchCountPerFilter[0]++
} else {
for (i in filters.indices) {
val limit = filters[i].limit
if ((limit == null || matchCountPerFilter[i] < limit) && filters[i].match(event)) {
matchCountPerFilter[i]++
}
}
}
}
@@ -155,15 +160,15 @@ suspend fun INostrClient.downloadFromRelay(
return totalEvents
}
suspend fun INostrClient.downloadFromRelay(
suspend fun INostrClient.reqBypassingRelayLimits(
relay: String,
baseFilters: List<Filter>,
filters: List<Filter>,
timeoutMs: Long = 30_000L,
onEvent: (Event) -> Unit,
): Int =
downloadFromRelay(
reqBypassingRelayLimits(
relay = RelayUrlNormalizer.normalize(relay),
baseFilters = baseFilters,
filters = filters,
timeoutMs = timeoutMs,
onEvent = onEvent,
)
@@ -21,11 +21,34 @@
package com.vitorpamplona.quartz.nip01Core.relay
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.BasicOkHttpWebSocket
import okhttp3.Interceptor
import okhttp3.OkHttpClient
import okhttp3.Request
import okhttp3.Response
class DefaultContentTypeInterceptor(
private val userAgentHeader: String,
) : Interceptor {
override fun intercept(chain: Interceptor.Chain): Response {
val originalRequest: Request = chain.request()
val requestWithUserAgent: Request =
originalRequest
.newBuilder()
.header("User-Agent", userAgentHeader)
.build()
return chain.proceed(requestWithUserAgent)
}
}
open class BaseNostrClientTest {
companion object {
val rootClient = OkHttpClient.Builder().build()
val rootClient =
OkHttpClient
.Builder()
.followRedirects(true)
.followSslRedirects(true)
.addInterceptor(DefaultContentTypeInterceptor("Amethyst/v1.05"))
.build()
val socketBuilder = BasicOkHttpWebSocket.Builder { url -> rootClient }
}
}
@@ -1,68 +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.quartz.nip01Core.relay
import com.vitorpamplona.quartz.nip01Core.metadata.MetadataEvent
import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.downloadFromRelay
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel
import kotlinx.coroutines.runBlocking
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertTrue
class NostrClientDownloadFromRelayTest : BaseNostrClientTest() {
@Test
fun testDownloadFromRelayReturnsMetadataEvents() =
runBlocking {
val appScope = CoroutineScope(Dispatchers.Default + SupervisorJob())
val client = NostrClient(socketBuilder, appScope)
val events = mutableListOf<com.vitorpamplona.quartz.nip01Core.core.Event>()
val totalFound =
client.downloadFromRelay(
relay = "wss://nos.lol",
baseFilters =
listOf(
Filter(
kinds = listOf(MetadataEvent.KIND),
limit = 10,
),
),
) { event ->
events.add(event)
}
client.disconnect()
appScope.cancel()
assertTrue(totalFound > 0, "Expected at least one event from wss://nos.lol")
assertTrue(events.isNotEmpty(), "Events list should not be empty")
events.forEach { event ->
assertEquals(MetadataEvent.KIND, event.kind, "All events should be kind ${MetadataEvent.KIND}")
}
}
}
@@ -0,0 +1,120 @@
/*
* Copyright (c) 2025 Vitor Pamplona
*
* Permission is hereby granted, free of charge, to any person obtaining a copy of
* this software and associated documentation files (the "Software"), to deal in
* the Software without restriction, including without limitation the rights to use,
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
* Software, and to permit persons to whom the Software is furnished to do so,
* subject to the following conditions:
*
* The above copyright notice and this permission notice shall be included in all
* copies or substantial portions of the Software.
*
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
*/
package com.vitorpamplona.quartz.nip01Core.relay
import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.metadata.MetadataEvent
import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.reqBypassingRelayLimits
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.nip02FollowList.ContactListEvent
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel
import kotlinx.coroutines.delay
import kotlinx.coroutines.runBlocking
import kotlin.test.Test
import kotlin.test.assertEquals
class NostrClientReqBypassingRelayLimitsTest : BaseNostrClientTest() {
@Test
fun testDownloadFromRelayReturnsMetadataEvents() =
runBlocking {
val appScope = CoroutineScope(Dispatchers.Default + SupervisorJob())
val client = NostrClient(socketBuilder, appScope)
val events = mutableListOf<Event>()
// nos.lol returns only 500 events per req
val totalFound =
client.reqBypassingRelayLimits(
relay = "wss://nos.lol",
filters =
listOf(
Filter(
kinds = listOf(MetadataEvent.KIND),
limit = 1000,
),
),
) { event ->
events.add(event)
}
client.disconnect()
delay(500)
appScope.cancel()
assertEquals(1000, totalFound, "Expected 1000 events from wss://nos.lol")
assertEquals(1000, events.size, "Events list should be 1000 events")
events.forEach { event ->
assertEquals(MetadataEvent.KIND, event.kind, "All events should be kind ${MetadataEvent.KIND}")
}
}
@Test
fun testDownloadFromRelayReturnsMetadataAndContactListEvents() =
runBlocking {
val appScope = CoroutineScope(Dispatchers.Default + SupervisorJob())
val client = NostrClient(socketBuilder, appScope)
val metadataEvents = mutableListOf<Event>()
val contactListEvents = mutableListOf<Event>()
// nos.lol returns only 500 events per req
val totalFound =
client.reqBypassingRelayLimits(
relay = "wss://nos.lol",
filters =
listOf(
Filter(
kinds = listOf(MetadataEvent.KIND),
limit = 1000,
),
Filter(
kinds = listOf(ContactListEvent.KIND),
limit = 1500,
),
),
) { event ->
if (event.kind == MetadataEvent.KIND) {
metadataEvents.add(event)
}
if (event.kind == ContactListEvent.KIND) {
contactListEvents.add(event)
}
}
client.disconnect()
delay(500)
appScope.cancel()
assertEquals(2500, totalFound, "Expected 1000 events from wss://nos.lol")
assertEquals(1000, metadataEvents.size, "Events list should be 1000 events")
assertEquals(1500, contactListEvents.size, "Events list should be 1000 events")
metadataEvents.forEach { event ->
assertEquals(MetadataEvent.KIND, event.kind, "All events should be kind ${MetadataEvent.KIND}")
}
contactListEvents.forEach { event ->
assertEquals(ContactListEvent.KIND, event.kind, "All events should be kind ${ContactListEvent.KIND}")
}
}
}