mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 11:48:24 +00:00
Merge pull request #4242 from vitorpamplona/fix/relay-filter-count-cap-reland
fix(relay): fit a REQ under a relay's filter-count cap by merging (re-land of #4239)
This commit is contained in:
+108
@@ -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<Filter>,
|
||||
): List<Filter> {
|
||||
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<Filter>,
|
||||
): List<Filter> {
|
||||
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<String, Unit>()
|
||||
|
||||
// The most filters a relay accepts in one REQ, learned from "invalid number of filters: N".
|
||||
private val maxFilters = ConcurrentMap<NormalizedRelayUrl, Int>()
|
||||
|
||||
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<Int> = 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<Filter>): List<Filter> {
|
||||
val out = mutableListOf<Filter>()
|
||||
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<Int> {
|
||||
val t = reason.lowercase()
|
||||
val kinds = mutableSetOf<Int>()
|
||||
|
||||
+169
@@ -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<Command>()
|
||||
|
||||
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<ReqCmd> {
|
||||
pool.onConnecting(url)
|
||||
val sent = mutableListOf<Command>()
|
||||
pool.syncState(url) { sent.add(it) }
|
||||
return sent.filterIsInstance<ReqCmd>()
|
||||
}
|
||||
|
||||
@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<ReqCmd>()
|
||||
.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)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user