diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayReqRefusals.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayReqRefusals.kt index 8aa828e211..3c8d87bca4 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayReqRefusals.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/RelayReqRefusals.kt @@ -22,6 +22,7 @@ package com.vitorpamplona.quartz.nip01Core.relay.client.pool import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl +import com.vitorpamplona.quartz.utils.Log import com.vitorpamplona.quartz.utils.concurrent.ConcurrentMap import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.StateFlow @@ -100,6 +101,7 @@ class RelayReqRefusals( reason: String, ): Boolean { learnDisallowedKinds(relay, reason) + learnMaxFilters(relay, reason) val candidate = classify(reason) ?: return false // NO_READS is the strictest verdict; once reached, nothing softens it. if (blocked[relay] == Policy.NO_READS) return false @@ -166,6 +168,23 @@ class RelayReqRefusals( fun narrow( relay: NormalizedRelayUrl, filters: List, + ): List { + val withoutRefusedKinds = stripRefusedKinds(relay, filters) + val cap = maxFilters[relay] ?: return withoutRefusedKinds + if (withoutRefusedKinds.size <= cap) return withoutRefusedKinds + val merged = mergeForCap(withoutRefusedKinds) + if (merged.size <= cap) return merged + // Still too many: a partial REQ the relay accepts beats a whole one it refuses. + val dropped = merged.drop(cap).map { it.kinds } + if (trimWarned.putIfAbsent("${relay.url} $dropped", Unit) == null) { + Log.w("RelayReqRefusals") { "${relay.url} caps REQs at $cap filters; sending $cap of ${merged.size}, dropping kinds $dropped" } + } + return merged.take(cap) + } + + private fun stripRefusedKinds( + relay: NormalizedRelayUrl, + filters: List, ): List { val refused = disallowedKinds[relay] ?: return filters return filters.mapNotNull { filter -> @@ -179,6 +198,20 @@ class RelayReqRefusals( } } + // Shapes already reported as trimmed, so a REQ re-decided on every EOSE warns once. + private val trimWarned = ConcurrentMap() + + // The most filters a relay accepts in one REQ, learned from "invalid number of filters: N". + private val maxFilters = ConcurrentMap() + + private fun learnMaxFilters( + relay: NormalizedRelayUrl, + reason: String, + ) { + val cap = parseMaxFilters(reason) ?: return + maxFilters.merge(relay, cap) { old, new -> minOf(old, new) } + } + fun disallowedKinds(relay: NormalizedRelayUrl): Set = disallowedKinds[relay] ?: emptySet() private fun classify(reason: String): Policy? { @@ -193,6 +226,81 @@ class RelayReqRefusals( private val KINDS_AFTER_MARKER = Regex("""kinds? (?:is |are )?not allowed:?\s*([0-9][0-9,\s]*)""") private val KIND_BEFORE_MARKER = Regex("""kind ([0-9]+) (?:is )?not allowed""") + // "invalid number of filters: 4" (strfry policy) — the relay refused N, so it takes fewer. + private val INVALID_FILTER_COUNT = Regex("""invalid number of filters:?\s*([0-9]+)""") + + fun parseMaxFilters(reason: String): Int? { + val refusedCount = + INVALID_FILTER_COUNT + .find(reason.lowercase()) + ?.groupValues + ?.get(1) + ?.toIntOrNull() ?: return null + return (refusedCount - 1).takeIf { it >= 1 } + } + + /** + * Merges filters that ask for the same thing except for ONE set of values (the kinds, + * one tag's values, the authors, or the ids) into a single filter with the union of those values + * and the earliest `since`. The result matches a superset of what each original + * matched, so nothing asked for is lost; an event the union adds is one some other + * filter in the same REQ already wanted, or a harmless duplicate. + * + * A filter with a `limit` is never merged: a limit per filter and a limit over the + * union are different requests. + */ + fun mergeForCap(filters: List): List { + val out = mutableListOf() + val remaining = filters.toMutableList() + while (remaining.isNotEmpty()) { + val head = remaining.removeAt(0) + if (head.limit != null) { + out.add(head) + continue + } + var merged = head + val iterator = remaining.iterator() + while (iterator.hasNext()) { + val candidate = iterator.next() + val union = unionIfMergeable(merged, candidate) ?: continue + merged = union + iterator.remove() + } + out.add(merged) + } + return out + } + + private fun unionIfMergeable( + a: Filter, + b: Filter, + ): Filter? { + if (b.limit != null) return null + if (a.until != b.until || a.search != b.search || a.tagsAll != b.tagsAll) return null + val sameKinds = a.kinds?.toSet() == b.kinds?.toSet() + val since = if (a.since == null || b.since == null) null else minOf(a.since, b.since) + + val sameIds = a.ids?.toSet() == b.ids?.toSet() + val sameAuthors = a.authors?.toSet() == b.authors?.toSet() + val aTags = a.tags ?: emptyMap() + val bTags = b.tags ?: emptyMap() + if (aTags.keys != bTags.keys) return null + val differingTags = aTags.keys.filter { aTags[it]?.toSet() != bTags[it]?.toSet() } + + val differences = (if (sameKinds) 0 else 1) + (if (sameIds) 0 else 1) + (if (sameAuthors) 0 else 1) + differingTags.size + return when { + differences == 0 -> a.copy(since = since) + differences > 1 -> null + !sameKinds -> if (a.kinds == null || b.kinds == null) null else a.copy(kinds = (a.kinds + b.kinds).distinct(), since = since) + !sameIds -> if (a.ids == null || b.ids == null) null else a.copy(ids = (a.ids + b.ids).distinct(), since = since) + !sameAuthors -> if (a.authors == null || b.authors == null) null else a.copy(authors = (a.authors + b.authors).distinct(), since = since) + else -> { + val key = differingTags.single() + a.copy(tags = aTags + (key to (aTags.getValue(key) + bTags.getValue(key)).distinct()), since = since) + } + } + } + fun parseDisallowedKinds(reason: String): Set { val t = reason.lowercase() val kinds = mutableSetOf() diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolRequestsFilterCapTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolRequestsFilterCapTest.kt new file mode 100644 index 0000000000..bc15867fcf --- /dev/null +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolRequestsFilterCapTest.kt @@ -0,0 +1,169 @@ +/* + * 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.client.pool + +import com.vitorpamplona.quartz.nip01Core.relay.client.single.IRelayClient +import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.ClosedMessage +import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.Command +import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.ReqCmd +import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertTrue + +/** + * relay.us.whitenoise.chat (strfry + policy) CLOSES any REQ with 4 or more filters + * (`ERROR: bad req: filter validation failed: invalid number of filters: N`) and + * advertises no max_filters in NIP-11. Marmot subscribes to group messages with one + * `{kinds:[445], #h:[group]}` filter per group in one REQ, so an account with 4+ + * groups on that relay loaded no group messages at all. + */ +class PoolRequestsFilterCapTest { + private val relay = NormalizedRelayUrl("wss://relay.us.whitenoise.chat/") + private val other = NormalizedRelayUrl("wss://nos.lol/") + + private class RecordingRelayClient( + override val url: NormalizedRelayUrl, + ) : IRelayClient { + val sent = mutableListOf() + + override fun connect() = Unit + + override fun needsToReconnect() = false + + override fun connectAndSyncFiltersIfDisconnected(ignoreRetryDelays: Boolean) = Unit + + override fun isConnected() = true + + override fun sendOrConnectAndSync(cmd: Command) { + sent.add(cmd) + } + + override fun sendIfConnected(cmd: Command) { + sent.add(cmd) + } + + override fun disconnect() = Unit + } + + private fun groupFilters(n: Int) = (1..n).map { Filter(kinds = listOf(445), tags = mapOf("h" to listOf("g$it")), since = 1000L + it) } + + private fun sync( + pool: PoolRequests, + url: NormalizedRelayUrl, + ): List { + pool.onConnecting(url) + val sent = mutableListOf() + pool.syncState(url) { sent.add(it) } + return sent.filterIsInstance() + } + + @Test + fun learnsTheCapFromTheRefusal() { + assertEquals(3, RelayReqRefusals.parseMaxFilters("ERROR: bad req: filter validation failed: invalid number of filters: 4")) + assertEquals(null, RelayReqRefusals.parseMaxFilters("ERROR: bad req: filter validation failed: kind not allowed: 21059")) + } + + @Test + fun perGroupFiltersAreMergedUnderTheCapAndSentAgainAtOnce() = + kotlinx.coroutines.test.runTest { + val pool = PoolRequests() + pool.addOrUpdate("groups", mapOf(relay to groupFilters(5)), null) + assertEquals(5, sync(pool, relay).single().filters.size) + + val client = RecordingRelayClient(relay) + pool.onIncomingMessage(client, ClosedMessage("groups", "ERROR: bad req: filter validation failed: invalid number of filters: 5")) + + val retry = + client.sent + .filterIsInstance() + .single() + .filters + assertEquals(1, retry.size, "five filters that differ only in #h become one") + assertEquals((1..5).map { "g$it" }.toSet(), retry.single().tags!!["h"]!!.toSet()) + assertEquals(listOf(445), retry.single().kinds) + assertEquals(1001L, retry.single().since, "the earliest since, so no group loses history") + } + + @Test + fun filtersThatCannotMergeAreTrimmedToTheCap() = + kotlinx.coroutines.test.runTest { + val pool = PoolRequests() + // Kinds AND authors differ, so no two of these can be merged. + val distinct = (1..5).map { Filter(kinds = listOf(it), authors = listOf("$it".repeat(64))) } + pool.addOrUpdate("mixed", mapOf(relay to distinct), null) + sync(pool, relay) + pool.onIncomingMessage(RecordingRelayClient(relay), ClosedMessage("mixed", "ERROR: bad req: filter validation failed: invalid number of filters: 4")) + + val filters = sync(pool, relay).single().filters + assertEquals(3, filters.size, "a partial REQ the relay accepts beats one it refuses") + } + + @Test + fun filtersWithALimitAreNotMerged() = + kotlinx.coroutines.test.runTest { + val pool = PoolRequests() + val limited = (1..5).map { Filter(kinds = listOf(1), authors = listOf("a$it"), limit = 20) } + pool.addOrUpdate("limited", mapOf(relay to limited), null) + sync(pool, relay) + pool.onIncomingMessage(RecordingRelayClient(relay), ClosedMessage("limited", "ERROR: bad req: filter validation failed: invalid number of filters: 4")) + + val filters = sync(pool, relay).single().filters + assertEquals(3, filters.size) + assertTrue(filters.all { it.authors!!.size == 1 }, "merging would change what each limit means") + } + + @Test + fun theCapOnlyAppliesToTheRelayThatRefused() = + kotlinx.coroutines.test.runTest { + val pool = PoolRequests() + pool.addOrUpdate("groups", mapOf(relay to groupFilters(5), other to groupFilters(5)), null) + sync(pool, relay) + pool.onIncomingMessage(RecordingRelayClient(relay), ClosedMessage("groups", "ERROR: bad req: filter validation failed: invalid number of filters: 5")) + + assertEquals(5, sync(pool, other).single().filters.size) + } + + @Test + fun anAccountsOwnListsMergeByKind() { + // Seen on device: the account's own lists (mute sets, bookmarks, pins, …) as + // several filters on the same author that differ only in their kinds. + val me = listOf("a".repeat(64)) + val lists = + listOf( + Filter(kinds = listOf(30000, 39089, 10000), authors = me), + Filter(kinds = listOf(10003, 30001, 30003), authors = me), + Filter(kinds = listOf(1984), authors = me), + ) + val merged = RelayReqRefusals.mergeForCap(lists).single() + assertEquals(setOf(30000, 39089, 10000, 10003, 30001, 30003, 1984), merged.kinds!!.toSet()) + assertEquals(me, merged.authors) + } + + @Test + fun twoDifferencesAtOnceAreNotMerged() { + // Different kinds AND different authors: the union would ask for pairs no filter wanted. + val a = Filter(kinds = listOf(1), authors = listOf("a".repeat(64))) + val b = Filter(kinds = listOf(7), authors = listOf("b".repeat(64))) + assertEquals(2, RelayReqRefusals.mergeForCap(listOf(a, b)).size) + } +}