From d6b8a54d8a7e6aa3f4800bbda2c62505e563b35e Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 6 Aug 2026 16:21:32 +0000 Subject: [PATCH 1/3] NEG-ERR: state the relay's max_sync_events on an overflow refusal MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A client that is refused for matching too much has exactly one thing to decide — how much smaller to ask next time — and no way to find out. NIP-11 has no field for max_sync_events, so the only route to a window the relay will answer is to guess and halve, and every wrong guess costs the relay the snapshot scan that produces the refusal. strfry already states the number in its rejection text; this makes it a first-class part of the frame. ["NEG-ERR", , ] unchanged, still what NIP-77 says ["NEG-ERR", , , ] when the refusal is about size Both mappers write the fourth element only when there is one, so a refusal with nothing to state is byte-identical to before, and both tolerate a non-numeric fourth element from someone else's relay. NegErrMessage.statedCap reads either form — the wire field or strfry's "(2431002 > 1000000)" prose — but only for a refusal that is about SIZE. That gate is the point of the property: a rate limit or a quota can carry numbers too, and it does not shrink when the window shrinks, so a client that mistook one for a cap would shrink its windows forever against a relay that has no size limit at all. The relay side sends its own configured cap for the same reason it is cheap: it had to know the number to refuse. --- .../kotlinSerialization/MessageKSerializer.kt | 9 ++ .../relay/server/NegSessionRegistry.kt | 8 +- .../quartz/nip77Negentropy/NegErrMessage.kt | 51 ++++++++++ .../nip77Negentropy/NegentropySettings.kt | 4 +- .../nip77Negentropy/NegErrMessageTest.kt | 96 +++++++++++++++++++ .../commands/toClient/MessageDeserializer.kt | 16 +++- .../commands/toClient/MessageSerializer.kt | 1 + .../nip77Negentropy/Nip77SerializationTest.kt | 64 +++++++++++++ 8 files changed, 243 insertions(+), 6 deletions(-) create mode 100644 quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegErrMessageTest.kt diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/kotlinSerialization/MessageKSerializer.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/kotlinSerialization/MessageKSerializer.kt index adba985212..93c0c360cb 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/kotlinSerialization/MessageKSerializer.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/kotlinSerialization/MessageKSerializer.kt @@ -46,6 +46,7 @@ import kotlinx.serialization.json.buildJsonArray import kotlinx.serialization.json.jsonArray import kotlinx.serialization.json.jsonObject import kotlinx.serialization.json.jsonPrimitive +import kotlinx.serialization.json.longOrNull object MessageKSerializer : KSerializer { override val descriptor: SerialDescriptor = @@ -112,6 +113,10 @@ object MessageKSerializer : KSerializer { is NegErrMessage -> { add(JsonPrimitive(value.subId)) add(JsonPrimitive(value.reason)) + // Only written when there is one: a three-element + // NEG-ERR is what NIP-77 describes, and that is what a + // refusal with nothing to state stays. + value.cap?.let { add(JsonPrimitive(it)) } } } } @@ -184,6 +189,10 @@ object MessageKSerializer : KSerializer { NegErrMessage( subId = array[1].jsonPrimitive.content, reason = if (array.size > 2) array[2].jsonPrimitive.content else "", + // Optional, and only a number: a relay that puts something + // else there is telling us nothing rather than breaking the + // frame. + cap = if (array.size > 3) array[3].jsonPrimitive.longOrNull else null, ) } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/NegSessionRegistry.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/NegSessionRegistry.kt index 66e1675e32..18b4a31925 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/NegSessionRegistry.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/NegSessionRegistry.kt @@ -104,7 +104,13 @@ class NegSessionRegistry( // `null` = matching set exceeds the cap (strfry-parity error). val sealedStorage = store.sealedNegentropyStorage(filters, maxEntries = settings.maxSyncEvents) if (sealedStorage == null) { - send(NegErrMessage(cmd.subId, "blocked: too many query results")) + // The cap rides along with the refusal. A client cannot discover + // this number any other way — NIP-11 has no field for it — so + // without it the only route to a window we WILL answer is guessing, + // halving, one refused NEG-OPEN at a time. Every one of those costs + // us the snapshot scan that produced this rejection, which makes + // stating it cheaper for the relay than staying quiet. + send(NegErrMessage(cmd.subId, "blocked: too many query results", settings.maxSyncEvents.toLong())) return } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegErrMessage.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegErrMessage.kt index a81e1ded1f..e191598d97 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegErrMessage.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegErrMessage.kt @@ -22,13 +22,64 @@ package com.vitorpamplona.quartz.nip77Negentropy import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.Message +/** + * `["NEG-ERR", , ]`, optionally followed by the relay's own + * `max_sync_events` when the refusal is about result-set size. + * + * That fourth element is not in NIP-77, but it is the only way a client learns + * the one number that decides how to ask again — no NIP-11 field carries it — + * and it is free for the relay to send, since it must know its own cap to have + * refused. strfry states it in the prose (`… too many records (2431002 > + * 1000000)`); [statedCap] reads either form. + * + * @property cap the fourth wire element, when present. + */ class NegErrMessage( val subId: String, val reason: String, + val cap: Long? = null, ) : Message { override fun label() = LABEL + /** + * The relay's negentropy cap if this refusal states one, from the wire + * field or from the prose, in that order. + * + * Only read for a refusal that is about SIZE ([isOverflow]). A quota or + * rate-limit refusal can carry numbers too, and sizing future windows + * against one of those would shrink every ask against a relay that has no + * size limit at all — while the limit that actually refused does not move + * however small the window gets. + */ + val statedCap: Long? + get() = if (!isOverflow(reason)) null else cap?.takeIf { it > 0 } ?: capInReason(reason) + companion object { const val LABEL = "NEG-ERR" + + /** `(2431002 > 1000000)` — the cap is the right-hand side. */ + private val COMPARISON = Regex("""\(\s*\d+\s*>\s*(\d+)\s*\)""") + + /** + * Does this reason mean "your query matched more than I will + * reconcile"? — as opposed to any other refusal, which no amount of + * window splitting will get past. + */ + fun isOverflow(reason: String): Boolean = + reason.contains("too many records", ignoreCase = true) || + reason.contains("too many results", ignoreCase = true) || + reason.contains("too many query results", ignoreCase = true) || + reason.contains("result set too large", ignoreCase = true) || + reason.contains("results too large", ignoreCase = true) || + reason.contains("max_sync_events", ignoreCase = true) + + /** The cap strfry writes into the refusal text, when it is there. */ + fun capInReason(reason: String): Long? = + COMPARISON + .find(reason) + ?.groupValues + ?.get(1) + ?.toLongOrNull() + ?.takeIf { it > 0 } } } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegentropySettings.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegentropySettings.kt index 896ec7c472..0bd1c1027d 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegentropySettings.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegentropySettings.kt @@ -32,7 +32,9 @@ package com.vitorpamplona.quartz.nip77Negentropy * unlimited). * @param maxSyncEvents Hard cap on the snapshot size for a single * NEG-OPEN. Mirrors strfry's `relay__negentropy__maxSyncEvents`. - * Overflow returns NEG-ERR `"blocked: too many query results"`. + * Overflow returns NEG-ERR `"blocked: too many query results"` + * carrying this number as its fourth element, so a client can size + * its next window instead of halving its way down to one. * @param maxSessionsPerConnection Cap on concurrent NEG sessions * held by one connection. strfry shares 200 with REQ subs; we * count NEG independently. Overflow sends NOTICE diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegErrMessageTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegErrMessageTest.kt new file mode 100644 index 0000000000..9adafce838 --- /dev/null +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegErrMessageTest.kt @@ -0,0 +1,96 @@ +/* + * 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.nip77Negentropy + +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertFalse +import kotlin.test.assertNull +import kotlin.test.assertTrue + +/** + * A stated cap is acted on — it sizes the next NEG-OPEN — so reading one out of + * a refusal that is not about size is worse than reading none at all: a quota or + * rate limit does not shrink when the window shrinks, so a client that mistook + * one for a cap would shrink its windows forever against a relay that has no + * size limit. + */ +class NegErrMessageTest { + @Test + fun capComesFromTheWireField() { + assertEquals(1_000_000L, NegErrMessage("s", "blocked: too many query results", 1_000_000L).statedCap) + } + + @Test + fun capComesFromStrfrysProseWhenTheFieldIsAbsent() { + val msg = NegErrMessage("s", "blocked: query matches too many records (2431002 > 1000000)") + assertEquals(1_000_000L, msg.statedCap) + } + + @Test + fun theWireFieldWinsOverTheProse() { + val msg = NegErrMessage("s", "blocked: too many records (5 > 10)", 1_000L) + assertEquals(1_000L, msg.statedCap) + } + + @Test + fun anOverflowWithNoNumberStatesNothing() { + assertNull(NegErrMessage("s", "blocked: too many query results").statedCap) + } + + @Test + fun aRateLimitIsNotACapHoweverManyNumbersItCarries() { + assertFalse(NegErrMessage.isOverflow("rate-limited: too many requests (30 > 10)")) + assertNull(NegErrMessage("s", "rate-limited: too many requests (30 > 10)", 10L).statedCap) + } + + @Test + fun refusalsThatAreNotAboutSizeStateNothing() { + listOf( + "auth-required: we only serve negentropy to authenticated users", + "blocked: pubkey is banned", + "error: negentropy disabled", + "closed: unknown subscription handle", + ).forEach { + assertFalse(NegErrMessage.isOverflow(it), "read as an overflow: $it") + assertNull(NegErrMessage("s", it, 42L).statedCap, "read a cap from: $it") + } + } + + @Test + fun theWordingsThatDoMeanOverflow() { + listOf( + "blocked: query matches too many records (5 > 1)", + "blocked: too many query results", + "error: result set too large", + "blocked: results too large", + "blocked: max_sync_events exceeded", + ).forEach { assertTrue(NegErrMessage.isOverflow(it), "not read as an overflow: $it") } + } + + @Test + fun aNonsensicalCapIsRefused() { + // Zero would wedge a client at a window that can never fit. + assertNull(NegErrMessage("s", "blocked: too many query results", 0L).statedCap) + assertNull(NegErrMessage("s", "blocked: too many records (5 > 0)").statedCap) + assertNull(NegErrMessage("s", "blocked: too many query results", -1L).statedCap) + } +} diff --git a/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/commands/toClient/MessageDeserializer.kt b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/commands/toClient/MessageDeserializer.kt index 7e45082421..a2d3e10df3 100644 --- a/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/commands/toClient/MessageDeserializer.kt +++ b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/commands/toClient/MessageDeserializer.kt @@ -121,10 +121,18 @@ class MessageDeserializer : StdDeserializer(Message::class.java) { } NegErrMessage.LABEL -> { - NegErrMessage( - subId = jp.nextTextValue(), - reason = jp.nextTextValue() ?: "", - ) + val subId = jp.nextTextValue() + val reason = jp.nextTextValue() ?: "" + // The optional fourth element, the relay's own cap. Read by + // stepping one token: anything that is not a number leaves + // the loop below to drain the frame, as before. + val cap = + if (jp.nextToken() == JsonToken.VALUE_NUMBER_INT) { + jp.longValue + } else { + null + } + NegErrMessage(subId, reason, cap) } else -> { diff --git a/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/commands/toClient/MessageSerializer.kt b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/commands/toClient/MessageSerializer.kt index 86274bf1e6..76671a396f 100644 --- a/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/commands/toClient/MessageSerializer.kt +++ b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/commands/toClient/MessageSerializer.kt @@ -129,6 +129,7 @@ class MessageSerializer : StdSerializer(Message::class.java) { is NegErrMessage -> { gen.writeString(msg.subId) gen.writeString(msg.reason) + msg.cap?.let { gen.writeNumber(it) } } } diff --git a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip77Negentropy/Nip77SerializationTest.kt b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip77Negentropy/Nip77SerializationTest.kt index 8b3bce89ba..7a71cc9a1b 100644 --- a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip77Negentropy/Nip77SerializationTest.kt +++ b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip77Negentropy/Nip77SerializationTest.kt @@ -124,6 +124,70 @@ class Nip77SerializationTest { assertEquals(msg.reason, jacksonDeserialized.reason) } + @Test + fun serializeNegErrMessageWithCap_matchesJackson() { + val msg = NegErrMessage("neg-sub1", "blocked: too many query results", 1_000_000L) + val jacksonJson = JacksonMapper.toJson(msg) + val kotlinJson = KotlinSerializationMapper.toJson(msg) + + assertEquals(jacksonJson, kotlinJson) + assertEquals("""["NEG-ERR","neg-sub1","blocked: too many query results",1000000]""", kotlinJson) + } + + @Test + fun serializeNegErrMessageWithoutCap_staysThreeElements() { + // NIP-77 describes a three-element NEG-ERR. A refusal with no cap to + // state must stay exactly that, rather than growing a fourth element + // every existing reader then has to tolerate. + val msg = NegErrMessage("neg-sub1", "closed: timeout") + assertEquals("""["NEG-ERR","neg-sub1","closed: timeout"]""", KotlinSerializationMapper.toJson(msg)) + assertEquals("""["NEG-ERR","neg-sub1","closed: timeout"]""", JacksonMapper.toJson(msg)) + } + + @Test + fun deserializeNegErrMessageWithCap_bothMappers() { + val json = """["NEG-ERR","neg-sub1","blocked: too many query results",1000000]""" + + val jackson = JacksonMapper.fromJsonToMessage(json) + assertTrue(jackson is NegErrMessage) + assertEquals(1_000_000L, jackson.cap) + assertEquals(1_000_000L, jackson.statedCap) + + val kotlin = KotlinSerializationMapper.fromJsonToMessage(json) + assertTrue(kotlin is NegErrMessage) + assertEquals(1_000_000L, kotlin.cap) + } + + @Test + fun deserializeNegErrMessageWithGarbageFourthElement_bothMappers() { + // A relay that puts something else there is telling us nothing; it must + // not break the frame that carries the reason. + val json = """["NEG-ERR","neg-sub1","blocked: too many query results","soon"]""" + + val jackson = JacksonMapper.fromJsonToMessage(json) + assertTrue(jackson is NegErrMessage) + assertEquals("blocked: too many query results", jackson.reason) + assertEquals(null, jackson.cap) + + val kotlin = KotlinSerializationMapper.fromJsonToMessage(json) + assertTrue(kotlin is NegErrMessage) + assertEquals("blocked: too many query results", kotlin.reason) + assertEquals(null, kotlin.cap) + } + + @Test + fun negErrMessageWithCap_crossDeserialization() { + val msg = NegErrMessage("neg-sub1", "blocked: too many records", 500_000L) + + val kotlinDeserialized = KotlinSerializationMapper.fromJsonToMessage(JacksonMapper.toJson(msg)) + assertTrue(kotlinDeserialized is NegErrMessage) + assertEquals(500_000L, kotlinDeserialized.cap) + + val jacksonDeserialized = JacksonMapper.fromJsonToMessage(KotlinSerializationMapper.toJson(msg)) + assertTrue(jacksonDeserialized is NegErrMessage) + assertEquals(500_000L, jacksonDeserialized.cap) + } + // ========================================================================= // NEG-OPEN Command (client-to-relay) Tests // ========================================================================= From 18998c897c655ad30195ff2282d656eb55e0e8d9 Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 6 Aug 2026 16:42:46 +0000 Subject: [PATCH 2/3] negentropy: size reconcile windows from the caller's own index MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A NEG-OPEN is all-or-nothing at both ends of the wire and neither end can see the other's size. The relay half has been handled since windowing landed — refuse, halve, retry. The client half has not: localEntries has to hold every matching (created_at, id) pair before the first NEG-OPEN goes out, so peak memory is a property of the CORPUS, not of the window. On a multi-million-event filter that list is the sync's high-water mark, and it is built even when the sync then splits into windows that each touch a fraction of it. NegentropyLocalIndex is that half. A caller whose store answers by range passes an index instead of a list, and the engine reads a window's worth at a time. count() is what makes it work: a window is sized BEFORE the round trip, so entriesFor() is only ever asked for something bounded. Callers that pass a list are unchanged — internally the list becomes an index that sorts once and binary-searches per window, exactly what the engine did inline before. targetWindow (0 = off, the old behaviour) turns the two signals into one loop. Our count splits a window before asking; their refusal shrinks the target — straight to the relay's stated cap where there is one, halved where there isn't — and windows that reconcile in one piece grow it back toward, never past, the caller's number. Neither side knows anything about the other and the same work queue absorbs both, which is what makes it adapt rather than need tuning. peerCap carries the relay's number back out, so a caller can persist it and start the NEXT sync at a window that fits. The local pre-split deliberately does NOT count against MAX_WINDOWS: that backstop exists for an overflow loop that never converges, while this split is driven by a number that provably halves with the range. Also here, because it is the same loop: page the window that overflowed rather than the whole filter. A second dense enough to exceed the cap is reachable — created_at has second granularity and is author-controlled — and negentropySyncOrFetch used to answer it by re-paging everything, including every window that had already reconciled cleanly. reconcileWindows now takes onUnreconcilableWindow and hands that window over; the sweep carries on with the rest of the range, so a dense second costs that second. Raw negentropySync/negentropyReconcile callers that pass no hook still get the exception, unchanged. pagedFallback stays conservative and now means "any part of this range came over REQ rather than a reconcile", with pagedWindows saying how much — the distinction matters to anyone recording coverage, since a paged walk booked as a completed reconcile would claim a range nothing compared. The existing over-cap test is updated rather than deleted: its ten events share one created_at, so the whole filter IS the un-reconcilable window — same events, now via the window path instead of by abandoning the sync. Its sibling test, that raw negentropySync still throws, is untouched. --- .../accessories/NegentropyLocalIndex.kt | 118 +++++++ .../accessories/NegentropySyncException.kt | 4 + .../NostrClientNegentropyFanOutExt.kt | 9 +- .../NostrClientNegentropySyncExt.kt | 330 +++++++++++++----- .../accessories/NegentropyLocalIndexTest.kt | 90 +++++ .../relay/NostrClientNegentropySyncTest.kt | 169 ++++++++- 6 files changed, 622 insertions(+), 98 deletions(-) create mode 100644 quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NegentropyLocalIndex.kt create mode 100644 quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NegentropyLocalIndexTest.kt diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NegentropyLocalIndex.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NegentropyLocalIndex.kt new file mode 100644 index 0000000000..6ed4d51619 --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NegentropyLocalIndex.kt @@ -0,0 +1,118 @@ +/* + * 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.accessories + +import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip01Core.store.IdAndTime + +/** + * The caller's own matching set, read one `created_at` window at a time. + * + * The list overloads of [negentropySync] / [negentropyReconcile] need every + * matching `(created_at, id)` pair before the first NEG-OPEN goes out, which + * makes peak memory a property of the corpus: a multi-million-event filter is + * a multi-million-entry list held for the whole sync, whether or not the sync + * ends up splitting into windows that each touch a fraction of it. + * + * A caller whose store can answer by range doesn't need that. Passing an index + * instead lets the window engine ask for a window's worth at a time, so the + * high-water mark becomes the size of one window — which the engine also sizes, + * from [count], before spending a round trip on it. + * + * Both methods are called on the reconciler coroutines, possibly concurrently + * when `reconcileConcurrency > 1`, and possibly more than once for the same + * window (a window that overflows is re-asked as halves). Implementations + * should be cheap and side-effect free; a store-backed one usually is, since + * these are index scans. + */ +interface NegentropyLocalIndex { + /** + * How many local events fall inside [window], or null when the store + * cannot answer cheaply. + * + * This is what lets the engine split a window BEFORE asking the relay for + * it — the only signal available about our own side, and the one that + * bounds what [entriesFor] will have to materialise. Null disables that + * pre-split for the window; the relay's own refusal is then the only thing + * that shrinks it, exactly as before this method existed. + */ + suspend fun count(window: Filter): Int? + + /** The `(created_at, id)` pairs inside [window]. Order does not matter. */ + suspend fun entriesFor(window: Filter): List + + companion object { + /** Nothing held locally: the sync downloads the relay's whole matched set. */ + val Empty: NegentropyLocalIndex = + object : NegentropyLocalIndex { + override suspend fun count(window: Filter) = 0 + + override suspend fun entriesFor(window: Filter) = emptyList() + } + + /** + * An index over a list already in memory — what the list overloads use, + * so they behave exactly as they did: sorted once, then binary-searched + * per window. + */ + fun of(entries: List): NegentropyLocalIndex = if (entries.isEmpty()) Empty else SortedListIndex(entries.sortedBy { it.createdAt }) + } +} + +private class SortedListIndex( + private val sorted: List, +) : NegentropyLocalIndex { + override suspend fun count(window: Filter): Int = slice(window).size + + override suspend fun entriesFor(window: Filter): List = slice(window) + + /** + * The `createdAt`-range slice of [sorted] (ascending by `createdAt`) that + * belongs to `[since, until]` (both inclusive, NIP-01 semantics). + * Binary-searched so window splits stay O(log n) over multi-million sets. + */ + private fun slice(window: Filter): List { + val since = window.since + val until = window.until + if (sorted.isEmpty() || (since == null && until == null)) return sorted + + val lo = since ?: 0L + val hi = until ?: Long.MAX_VALUE + + // first index with createdAt >= lo + var start = 0 + var e = sorted.size + while (start < e) { + val mid = (start + e) ushr 1 + if (sorted[mid].createdAt < lo) start = mid + 1 else e = mid + } + + // first index with createdAt > hi + var end = start + e = sorted.size + while (end < e) { + val mid = (end + e) ushr 1 + if (sorted[mid].createdAt <= hi) end = mid + 1 else e = mid + } + + return if (start >= end) emptyList() else sorted.subList(start, end) + } +} diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NegentropySyncException.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NegentropySyncException.kt index f6018a8098..1a6dacc98d 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NegentropySyncException.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NegentropySyncException.kt @@ -43,12 +43,16 @@ import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl * @property window the filter slice that failed. * @property reason machine-readable category — branch on this to recover. * @property detail the underlying specifics (a relay's `NEG-ERR` text, `timeout`, …). + * @property cap the relay's own `max_sync_events` when its refusal stated one + * (see [com.vitorpamplona.quartz.nip77Negentropy.NegErrMessage.statedCap]). + * Worth persisting per relay: it is what sizes the first window next time. */ class NegentropySyncException( val relay: NormalizedRelayUrl, val window: Filter, val reason: Reason, val detail: String, + val cap: Long? = null, ) : Exception("NIP-77 sync of $relay failed ($reason): $detail") { enum class Reason { /** diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropyFanOutExt.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropyFanOutExt.kt index 7ac08647f5..e941af4417 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropyFanOutExt.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropyFanOutExt.kt @@ -77,6 +77,8 @@ suspend fun negentropySyncFanOut( relay: NormalizedRelayUrl, filter: Filter, localEntries: List = emptyList(), + localIndex: NegentropyLocalIndex? = null, + targetWindow: Int = 0, maxEvents: Int = 0, reqsPerClient: Int = 10, fetchBatch: Int = 250, @@ -135,8 +137,6 @@ suspend fun negentropySyncFanOut( val producer = launch { try { - val sorted = - if (localEntries.size > 1) localEntries.sortedBy { it.createdAt } else localEntries // Windows reconcile round-robin ACROSS the clients so // server-side snapshot builds parallelize per connection // (a single connection produced ids at only ~9k/s and @@ -145,17 +145,18 @@ suspend fun negentropySyncFanOut( clients = clients, relay = relay, filter = filter, - localEntries = sorted, + local = localIndex ?: NegentropyLocalIndex.of(localEntries), idleTimeoutMs = idleTimeoutMs, batchSize = fetchBatch, reconcileConcurrency = reconcileConcurrency, + targetWindow = targetWindow, onWindow = { windows.incrementAndFetch() }, onNeed = { need.addAndFetch(it) }, onHave = { have.addAndFetch(it) }, sendNeedBatch = { batch -> idBatches.send(batch) }, - sendHaveBatch = if (localEntries.isEmpty()) null else { _ -> }, + sendHaveBatch = if (localEntries.isEmpty() && localIndex == null) null else { _ -> }, ) } finally { idBatches.close() diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncExt.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncExt.kt index 611720876a..6595cc763e 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncExt.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncExt.kt @@ -66,12 +66,17 @@ import kotlin.math.min * @property downloaded distinct events actually delivered through `onEvent`. * @property windows number of `created_at` windows the matched set was split * into (`1` when the relay reconciled the whole filter in one shot). + * @property peerCap the relay's own `max_sync_events`, when a refusal during + * this sync stated one. Worth persisting per relay: it is the number that + * sizes the first window of the NEXT sync, and it is not discoverable any + * other way. */ class NegentropySyncResult( val needCount: Int, val haveCount: Int, val downloaded: Int, val windows: Int, + val peerCap: Long? = null, ) /** @@ -162,12 +167,16 @@ suspend fun INostrClient.negentropySync( reconcileConcurrency: Int = 1, idBufferBatches: Int = maxConcurrentReqs * 4, localEntries: List = emptyList(), + localIndex: NegentropyLocalIndex? = null, + targetWindow: Int = 0, + onUnreconcilableWindow: (suspend (Filter) -> Unit)? = null, onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null, onEvent: suspend (Event) -> Unit, ): NegentropySyncResult { val need = AtomicInt(0) val windows = AtomicInt(0) var downloaded = 0 + var peerCap: Long? = null // Pin the relay in the pool's "desired" set for the whole sync. A NEG-OPEN is not // a REQ, so during a reconcile round (before that window's first download REQ @@ -195,8 +204,11 @@ suspend fun INostrClient.negentropySync( maxConcurrentReqs = maxConcurrentReqs, reconcileConcurrency = reconcileConcurrency, idBufferBatches = idBufferBatches, - localEntries = localEntries, + local = localIndex ?: NegentropyLocalIndex.of(localEntries), + targetWindow = targetWindow, + onUnreconcilableWindow = onUnreconcilableWindow, onWindow = { windows.incrementAndFetch() }, + onPeerCap = { peerCap = it }, // Only accumulate here; progress is reported from the // single consumer loop below so the user callback is never // invoked from two coroutines at once. @@ -228,6 +240,7 @@ suspend fun INostrClient.negentropySync( haveCount = 0, downloaded = downloaded, windows = windows.load(), + peerCap = peerCap, ) } @@ -241,6 +254,9 @@ suspend fun INostrClient.negentropySync( reconcileConcurrency: Int = 1, idBufferBatches: Int = maxConcurrentReqs * 4, localEntries: List = emptyList(), + localIndex: NegentropyLocalIndex? = null, + targetWindow: Int = 0, + onUnreconcilableWindow: (suspend (Filter) -> Unit)? = null, onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null, onEvent: suspend (Event) -> Unit, ): NegentropySyncResult = @@ -254,6 +270,9 @@ suspend fun INostrClient.negentropySync( reconcileConcurrency = reconcileConcurrency, idBufferBatches = idBufferBatches, localEntries = localEntries, + localIndex = localIndex, + targetWindow = targetWindow, + onUnreconcilableWindow = onUnreconcilableWindow, onProgress = onProgress, onEvent = onEvent, ) @@ -263,16 +282,28 @@ suspend fun INostrClient.negentropySync( * * @property downloaded distinct events delivered through `onEvent` (across whichever * path ran). - * @property pagedFallback `true` if negentropy could not reconcile and the events - * came from [fetchAllPages] instead. + * @property pagedFallback `true` if ANY part of the range came from + * [fetchAllPages] rather than a reconcile — either the whole filter (the + * relay could not reconcile at all) or the individual windows counted by + * [pagedWindows]. Deliberately conservative: a caller recording what it has + * covered must not book a paged walk as a completed reconcile, and one + * un-reconcilable second in the range is enough to make that claim untrue. * @property negentropy the negentropy outcome when it succeeded; `null` on fallback. - * @property fallbackCause why negentropy was abandoned; `null` when it succeeded. + * @property fallbackCause why negentropy was abandoned for the WHOLE filter; + * `null` when it was not — including when individual windows were paged, which + * have no single cause between them. + * @property pagedWindows how many individual `created_at` windows were paged + * inside an otherwise-successful negentropy sync — seconds so dense the relay + * would not reconcile them at any window size. `0` for almost every sync; + * non-zero means part of the range came over REQ and is subject to a paged + * walk's limits rather than a reconcile's guarantees. */ class NegentropyOrFetchResult( val downloaded: Int, val pagedFallback: Boolean, val negentropy: NegentropySyncResult?, val fallbackCause: NegentropySyncException?, + val pagedWindows: Int = 0, ) /** @@ -308,11 +339,14 @@ suspend fun INostrClient.negentropySyncOrFetch( reconcileConcurrency: Int = 1, idBufferBatches: Int = maxConcurrentReqs * 4, localEntries: List = emptyList(), + localIndex: NegentropyLocalIndex? = null, + targetWindow: Int = 0, onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null, onEvent: suspend (Event) -> Unit, ): NegentropyOrFetchResult { val seen = HashSet() var delivered = 0 + var pagedWindows = 0 // Shared dedup + cap across both phases. Returns true if the event was new and // delivered. Both phases run sequentially, so no concurrent access. @@ -337,9 +371,32 @@ suspend fun INostrClient.negentropySyncOrFetch( reconcileConcurrency = reconcileConcurrency, idBufferBatches = idBufferBatches, localEntries = localEntries, + localIndex = localIndex, + targetWindow = targetWindow, + // One second the relay will not reconcile at any size costs + // that second, not the sync. Without this the exception below + // catches it and re-pages the WHOLE filter — every window that + // already reconciled cleanly walked again over REQ, which on a + // large corpus is the entire cost negentropy was there to save. + onUnreconcilableWindow = { window -> + pagedWindows++ + val pageTimeoutMs = if (idleTimeoutMs > 0) idleTimeoutMs else DEFAULT_DOWNLOAD_IDLE_MS + fetchAllPages(relay, listOf(window), pageTimeoutMs) { event -> + if (accept(event)) onProgress?.invoke(delivered, delivered) + } + }, onProgress = onProgress, ) { accept(it) } - NegentropyOrFetchResult(delivered, pagedFallback = false, negentropy = result, fallbackCause = null) + NegentropyOrFetchResult( + delivered, + // Any paged window makes this not a clean reconcile — see the + // property doc: under-reporting it would let a caller record + // coverage it never compared. + pagedFallback = pagedWindows > 0, + negentropy = result, + fallbackCause = null, + pagedWindows = pagedWindows, + ) } catch (e: NegentropySyncException) { // Negentropy couldn't enumerate the set — page the whole filter instead, // skipping anything the negentropy attempt already delivered. fetchAllPages @@ -349,7 +406,13 @@ suspend fun INostrClient.negentropySyncOrFetch( fetchAllPages(relay, listOf(pageFilter), pageTimeoutMs) { event -> if (accept(event)) onProgress?.invoke(delivered, delivered) } - NegentropyOrFetchResult(delivered, pagedFallback = true, negentropy = null, fallbackCause = e) + NegentropyOrFetchResult( + delivered, + pagedFallback = true, + negentropy = null, + fallbackCause = e, + pagedWindows = pagedWindows, + ) } } @@ -363,6 +426,8 @@ suspend fun INostrClient.negentropySyncOrFetch( reconcileConcurrency: Int = 1, idBufferBatches: Int = maxConcurrentReqs * 4, localEntries: List = emptyList(), + localIndex: NegentropyLocalIndex? = null, + targetWindow: Int = 0, onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null, onEvent: suspend (Event) -> Unit, ): NegentropyOrFetchResult = @@ -376,6 +441,8 @@ suspend fun INostrClient.negentropySyncOrFetch( reconcileConcurrency = reconcileConcurrency, idBufferBatches = idBufferBatches, localEntries = localEntries, + localIndex = localIndex, + targetWindow = targetWindow, onProgress = onProgress, onEvent = onEvent, ) @@ -411,9 +478,12 @@ private suspend fun INostrClient.syncPipeline( maxConcurrentReqs: Int, reconcileConcurrency: Int, idBufferBatches: Int, - localEntries: List, + local: NegentropyLocalIndex, + targetWindow: Int, onWindow: () -> Unit, onNeed: (Int) -> Unit, + onPeerCap: ((Long) -> Unit)?, + onUnreconcilableWindow: (suspend (Filter) -> Unit)?, deliver: suspend (Event) -> Unit, ) = coroutineScope { val idBatches = Channel>(idBufferBatches.coerceAtLeast(1)) @@ -430,21 +500,20 @@ private suspend fun INostrClient.syncPipeline( } } - // reconcileWindows needs the local set sorted by createdAt (it binary-searches - // each window's slice). Empty/singleton sets are already trivially sorted. - val sortedLocal = if (localEntries.size > 1) localEntries.sortedBy { it.createdAt } else localEntries - reconcileWindows( clients = listOf(this@syncPipeline), relay = relay, filter = filter, - localEntries = sortedLocal, + local = local, idleTimeoutMs = idleTimeoutMs, batchSize = fetchBatch, reconcileConcurrency = reconcileConcurrency, + targetWindow = targetWindow, onWindow = onWindow, onNeed = onNeed, onHave = {}, + onPeerCap = onPeerCap, + onUnreconcilableWindow = onUnreconcilableWindow, sendNeedBatch = { batch -> idBatches.send(batch) }, sendHaveBatch = null, ) @@ -455,16 +524,29 @@ private suspend fun INostrClient.syncPipeline( /** * The shared window engine behind [negentropySync] and [negentropyReconcile]: - * reconciles [filter] against [localEntries], splitting into `created_at` - * windows whenever the relay rejects the set as too large, with up to - * [reconcileConcurrency] windows reconciling at once from a shared work - * queue. Each window's local subset is sliced out of [localEntries] (which - * MUST be sorted by `createdAt`) so both sides always reconcile the same - * slice of the timeline. + * reconciles [filter] against [local], splitting into `created_at` windows, + * with up to [reconcileConcurrency] windows reconciling at once from a shared + * work queue. Each window reconciles against that window's slice of [local], so + * both sides always compare the same slice of the timeline. + * + * Two independent things split a window, and the same queue absorbs both: + * + * - **The relay refuses it** (strfry's `max_sync_events`). Known only after a + * round trip, and the only signal available about THEIR size. + * - **We hold more than [targetWindow] in it**, per [NegentropyLocalIndex.count], + * which is known before the round trip and is what bounds the entries this + * engine asks [local] to materialise. Off when [targetWindow] is `0` (the + * default), which is the pre-existing behaviour: one window until refused. + * + * Neither side can see the other's size, so [targetWindow] adapts within the + * sync: a refusal shrinks it — straight to the relay's own cap when the refusal + * states one ([NegErrMessage.statedCap]), halved when it does not — and windows + * that reconcile in one piece grow it back toward, never past, the caller's + * number. * * Throws [NegentropySyncException] for any window negentropy cannot reconcile - * (a minimal window still over the cap, or an unavailable/erroring relay); the - * failure cancels the whole scope. + * (a minimal window still over the cap with no [onUnreconcilableWindow] to hand + * it to, or an unavailable/erroring relay); the failure cancels the whole scope. */ @OptIn(ExperimentalAtomicApi::class) internal suspend fun reconcileWindows( @@ -474,13 +556,19 @@ internal suspend fun reconcileWindows( clients: List, relay: NormalizedRelayUrl, filter: Filter, - localEntries: List, + local: NegentropyLocalIndex, idleTimeoutMs: Long, batchSize: Int, reconcileConcurrency: Int, + targetWindow: Int = 0, onWindow: () -> Unit, onNeed: (Int) -> Unit, onHave: (Int) -> Unit, + onPeerCap: ((Long) -> Unit)? = null, + // Given a minimal window the relay will not reconcile at any size, instead + // of throwing. The caller drains it however it can (paging it over REQ) and + // the sweep carries on with the rest of the filter. + onUnreconcilableWindow: (suspend (Filter) -> Unit)? = null, sendNeedBatch: suspend (List) -> Unit, sendHaveBatch: (suspend (List) -> Unit)?, ) = coroutineScope { @@ -503,6 +591,33 @@ internal suspend fun reconcileWindows( // the tens–hundreds, so the cap is orders of magnitude above any real sync. val totalWindows = AtomicInt(1) + // The largest window this sync will ask for, in events. Shrinks on a + // refusal, recovers toward the caller's number on clean windows, and is + // read only where a local count exists to compare it against — with + // targetWindow at 0 nothing below this line does anything. + val budget = AtomicInt(targetWindow) + + // Splits a window in two and queues both halves. Returns false when the + // window is already minimal — `created_at` is in seconds, so that is the + // floor, not a tuning choice. + fun splitInto( + pendingWindow: Filter, + lo: Long, + hi: Long, + ): Boolean { + if (hi - lo <= MIN_WINDOW_SECONDS) return false + val mid = lo + (hi - lo) / 2 + remaining.incrementAndFetch() + // The lower child gets the finite midpoint; the upper child KEEPS this + // window's original `until` (which may be null = unbounded). Replacing + // null with `now()` here would drop every event dated after now() + // (clock skew) once any split happens, while the un-split path would + // have included them. + pending.trySend(pendingWindow.copy(since = lo, until = mid)) + pending.trySend(pendingWindow.copy(since = mid + 1, until = pendingWindow.until)) + return true + } + val reconcilers = List(reconcileConcurrency.coerceAtLeast(1)) { reconcilerIndex -> launch { @@ -510,11 +625,27 @@ internal suspend fun reconcileWindows( for (window in pending) { coroutineContext.ensureActive() + val lo = window.since ?: 0L + val hi = window.until ?: TimeUtils.now() + + // Our own side, before the round trip. Deliberately NOT + // counted against MAX_WINDOWS: that backstop guards against + // an overflow loop that never converges, while this split is + // driven by a number that provably halves with the range. + val ceiling = budget.load() + if (ceiling > 0 && hi - lo > MIN_WINDOW_SECONDS) { + val mine = local.count(window) + if (mine != null && mine > ceiling) { + splitInto(window, lo, hi) + continue + } + } + val outcome = client.reconcileStreaming( relay = relay, filter = window, - localEntries = entriesForWindow(localEntries, window.since, window.until), + localEntries = local.entriesFor(window), idleTimeoutMs = idleTimeoutMs, fetchBatch = batchSize, onNeed = onNeed, @@ -526,21 +657,53 @@ internal suspend fun reconcileWindows( when (outcome) { is ReconcileOutcome.Complete -> { onWindow() + // A window that fitted is evidence the budget can + // recover — gently, and never past what the caller + // asked for, so a sync that met one dense stretch + // does not stay small for the rest of the timeline. + if (targetWindow > 0) { + val now = budget.load() + if (now < targetWindow) { + budget.store(minOf(targetWindow, (now * BUDGET_GROWTH).toInt().coerceAtLeast(now + 1))) + } + } if (remaining.decrementAndFetch() == 0) pending.close() } is ReconcileOutcome.Overflow -> { - val lo = window.since ?: 0L - val hi = window.until ?: TimeUtils.now() + // What they will take, when they said so: one step + // instead of a halving ladder, for this sync and — + // via onPeerCap — for whatever the caller persists. + outcome.cap?.let { cap -> + onPeerCap?.invoke(cap) + if (targetWindow > 0) { + val fitted = (cap * CAP_MARGIN).toInt().coerceAtLeast(1) + if (fitted < budget.load()) budget.store(fitted) + } + } + if (outcome.cap == null && targetWindow > 0) { + // No number to go on: halve and find out. + budget.store((budget.load() / 2).coerceAtLeast(1)) + } if (hi - lo <= MIN_WINDOW_SECONDS) { - // A minimal window that still overflows: negentropy - // genuinely can't enumerate this slice. Surface it — - // paging is the caller's call. + // A minimal window that still overflows: + // negentropy genuinely can't enumerate this + // slice. Hand it to the caller if it has a way + // to drain it, otherwise surface it — paging is + // the caller's call either way. + val fallback = onUnreconcilableWindow + if (fallback != null) { + fallback(window) + onWindow() + if (remaining.decrementAndFetch() == 0) pending.close() + continue + } throw NegentropySyncException( relay = relay, window = window, reason = NegentropySyncException.Reason.OVER_MAX_SYNC_EVENTS, detail = "created_at window [$lo, $hi] still exceeds the relay's max_sync_events", + cap = outcome.cap, ) } if (totalWindows.addAndFetch(2) > MAX_WINDOWS) { @@ -554,15 +717,7 @@ internal suspend fun reconcileWindows( detail = "created_at window split exceeded $MAX_WINDOWS windows without converging; the relay likely rejects negentropy with an overflow-looking error", ) } - val mid = lo + (hi - lo) / 2 - remaining.incrementAndFetch() - // The lower child gets the finite midpoint; the upper child - // KEEPS this window's original `until` (which may be null = - // unbounded). Replacing null with `now()` here would drop - // every event dated after now() (clock skew) once any split - // happens, while the un-split path would have included them. - pending.send(window.copy(since = lo, until = mid)) - pending.send(window.copy(since = mid + 1, until = window.until)) + splitInto(window, lo, hi) } is ReconcileOutcome.Failed -> @@ -580,46 +735,17 @@ internal suspend fun reconcileWindows( reconcilers.joinAll() } -/** - * The `createdAt`-range slice of [sorted] (ascending by `createdAt`) that - * belongs to the window `[since, until]` (both inclusive, NIP-01 semantics). - * Binary-searched so window splits stay O(log n) over multi-million local sets. - */ -private fun entriesForWindow( - sorted: List, - since: Long?, - until: Long?, -): List { - if (sorted.isEmpty() || (since == null && until == null)) return sorted - - val lo = since ?: 0L - val hi = until ?: Long.MAX_VALUE - - // first index with createdAt >= lo - var start = 0 - var e = sorted.size - while (start < e) { - val mid = (start + e) ushr 1 - if (sorted[mid].createdAt < lo) start = mid + 1 else e = mid - } - - // first index with createdAt > hi - var end = start - e = sorted.size - while (end < e) { - val mid = (end + e) ushr 1 - if (sorted[mid].createdAt <= hi) end = mid + 1 else e = mid - } - - return if (start >= end) emptyList() else sorted.subList(start, end) -} - private sealed interface ReconcileOutcome { /** Reconciliation completed; every id was streamed to the downloader. */ object Complete : ReconcileOutcome - /** Relay rejected the set as too large (strfry `max_sync_events`). */ - object Overflow : ReconcileOutcome + /** + * Relay rejected the set as too large (strfry `max_sync_events`). + * [cap] is the relay's own limit when the refusal stated one. + */ + class Overflow( + val cap: Long?, + ) : ReconcileOutcome /** Reconciliation could not complete; [detail] says why. */ class Failed( @@ -633,11 +759,14 @@ private sealed interface ReconcileOutcome { * @property needCount ids the relay has that the local set lacks (streamed to `onNeedIds`). * @property haveCount ids the local set has that the relay lacks (streamed to `onHaveIds`). * @property windows number of `created_at` windows the reconcile split into. + * @property peerCap the relay's own `max_sync_events`, when a refusal during + * this reconcile stated one. */ class NegentropyReconcileResult( val needCount: Int, val haveCount: Int, val windows: Int, + val peerCap: Long? = null, ) /** @@ -682,15 +811,19 @@ suspend fun INostrClient.negentropyReconcile( relay: NormalizedRelayUrl, filter: Filter, localEntries: List = emptyList(), + localIndex: NegentropyLocalIndex? = null, + targetWindow: Int = 0, batchSize: Int = 500, idleTimeoutMs: Long = 120_000L, reconcileConcurrency: Int = 1, + onUnreconcilableWindow: (suspend (Filter) -> Unit)? = null, onHaveIds: (suspend (List) -> Unit)? = null, onNeedIds: suspend (List) -> Unit, ): NegentropyReconcileResult { val need = AtomicInt(0) val have = AtomicInt(0) val windows = AtomicInt(0) + var peerCap: Long? = null // Same connection-pinning trick as negentropySync: a NEG-OPEN is not a REQ, // so without a live subscription the pool would consider the relay unwanted @@ -698,24 +831,20 @@ suspend fun INostrClient.negentropyReconcile( val keepAliveSubId = newSubId() subscribe(keepAliveSubId, mapOf(relay to listOf(Filter(ids = listOf(KEEP_ALIVE_ID)))), null) try { - val sorted = - if (localEntries.size > 1) { - localEntries.sortedBy { it.createdAt } - } else { - localEntries - } - reconcileWindows( clients = listOf(this), relay = relay, filter = filter, - localEntries = sorted, + local = localIndex ?: NegentropyLocalIndex.of(localEntries), idleTimeoutMs = idleTimeoutMs, batchSize = batchSize, reconcileConcurrency = reconcileConcurrency, + targetWindow = targetWindow, onWindow = { windows.incrementAndFetch() }, onNeed = { need.addAndFetch(it) }, onHave = { have.addAndFetch(it) }, + onPeerCap = { peerCap = it }, + onUnreconcilableWindow = onUnreconcilableWindow, sendNeedBatch = onNeedIds, sendHaveBatch = onHaveIds, ) @@ -727,6 +856,7 @@ suspend fun INostrClient.negentropyReconcile( needCount = need.load(), haveCount = have.load(), windows = windows.load(), + peerCap = peerCap, ) } @@ -734,9 +864,12 @@ suspend fun INostrClient.negentropyReconcile( relay: String, filter: Filter, localEntries: List = emptyList(), + localIndex: NegentropyLocalIndex? = null, + targetWindow: Int = 0, batchSize: Int = 500, idleTimeoutMs: Long = 120_000L, reconcileConcurrency: Int = 1, + onUnreconcilableWindow: (suspend (Filter) -> Unit)? = null, onHaveIds: (suspend (List) -> Unit)? = null, onNeedIds: suspend (List) -> Unit, ): NegentropyReconcileResult = @@ -744,9 +877,12 @@ suspend fun INostrClient.negentropyReconcile( relay = RelayUrlNormalizer.normalize(relay), filter = filter, localEntries = localEntries, + localIndex = localIndex, + targetWindow = targetWindow, batchSize = batchSize, idleTimeoutMs = idleTimeoutMs, reconcileConcurrency = reconcileConcurrency, + onUnreconcilableWindow = onUnreconcilableWindow, onHaveIds = onHaveIds, onNeedIds = onNeedIds, ) @@ -895,7 +1031,7 @@ private suspend fun INostrClient.reconcileStreaming( clock.bump() if (msg.subId == subId) { sawNegFrame = true - incoming.trySend(NegFrame.Err(msg.reason)) + incoming.trySend(NegFrame.Err(msg.reason, msg.statedCap)) } } @@ -965,7 +1101,11 @@ private suspend fun INostrClient.reconcileStreaming( when (frame) { is NegFrame.Err -> - return if (isOverflow(frame.reason)) ReconcileOutcome.Overflow else ReconcileOutcome.Failed(frame.reason) + return if (isOverflow(frame.reason)) { + ReconcileOutcome.Overflow(frame.cap) + } else { + ReconcileOutcome.Failed(frame.reason) + } is NegFrame.Msg -> { val result = session.processMessage(frame.payload) @@ -1014,6 +1154,8 @@ private sealed interface NegFrame { class Err( val reason: String, + // The relay's own max_sync_events, when the refusal stated one. + val cap: Long? = null, ) : NegFrame } @@ -1041,13 +1183,7 @@ private sealed interface NegFrame { * [reconcileWindows] also caps the total window count as a wording-independent * backstop, so a novel overflow-looking-but-not-shrinking error can never storm. */ -internal fun isOverflow(reason: String): Boolean = - reason.contains("too many records", ignoreCase = true) || - reason.contains("too many results", ignoreCase = true) || - reason.contains("too many query results", ignoreCase = true) || - reason.contains("result set too large", ignoreCase = true) || - reason.contains("results too large", ignoreCase = true) || - reason.contains("max_sync_events", ignoreCase = true) +internal fun isOverflow(reason: String): Boolean = NegErrMessage.isOverflow(reason) /** * A relay that advertises NIP-77 but refuses it at runtime signals the refusal with @@ -1159,6 +1295,22 @@ private const val MIN_WINDOW_SECONDS = 1L */ private const val MAX_WINDOWS = 100_000 +/** + * How much of a relay's stated `max_sync_events` a window actually aims for. + * The margin absorbs what the relay gains between stating that number and + * answering the next NEG-OPEN — asking for exactly the cap would be refused + * again by anything still being written to. + */ +private const val CAP_MARGIN = 0.8 + +/** + * How fast a shrunk window grows back toward the caller's target, per window + * that reconciled in one piece. Multiplicative and gentle on purpose: too small + * costs an extra round trip, too big costs a refused NEG-OPEN plus the snapshot + * scan the relay did before refusing it. + */ +private const val BUDGET_GROWTH = 1.25 + /** Bounded buffer between the download workers and the single delivery consumer. */ private const val DELIVERY_BUFFER = 256 diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NegentropyLocalIndexTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NegentropyLocalIndexTest.kt new file mode 100644 index 0000000000..7967d1695d --- /dev/null +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NegentropyLocalIndexTest.kt @@ -0,0 +1,90 @@ +/* + * 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.accessories + +import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip01Core.store.IdAndTime +import kotlinx.coroutines.test.runTest +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertTrue + +/** + * The list-backed index is what the `localEntries` overloads become, so its + * slicing has to keep NIP-01's inclusive `since`/`until` exactly: a window that + * dropped its boundary second would leave events neither side ever compares, + * and the two sides of a reconcile would disagree about what the window holds. + */ +class NegentropyLocalIndexTest { + private fun idAt(second: Long) = IdAndTime(second, second.toString().padStart(64, '0')) + + private val index = NegentropyLocalIndex.of((1000L..1009L).map { idAt(it) }) + + private fun window( + since: Long?, + until: Long?, + ) = Filter(kinds = listOf(1), since = since, until = until) + + @Test + fun bothBoundsAreInclusive() = + runTest { + assertEquals(3, index.count(window(1002, 1004))) + assertEquals(listOf(1002L, 1003L, 1004L), index.entriesFor(window(1002, 1004)).map { it.createdAt }) + } + + @Test + fun anUnboundedSideReachesTheEnd() = + runTest { + assertEquals(5, index.count(window(1005, null))) + assertEquals(6, index.count(window(null, 1005))) + assertEquals(10, index.count(window(null, null))) + } + + @Test + fun aWindowOutsideEverythingIsEmpty() = + runTest { + assertEquals(0, index.count(window(2000, 3000))) + assertTrue(index.entriesFor(window(2000, 3000)).isEmpty()) + } + + @Test + fun aSingleSecondWindowHoldsThatSecond() = + runTest { + assertEquals(1, index.count(window(1007, 1007))) + assertEquals(listOf(1007L), index.entriesFor(window(1007, 1007)).map { it.createdAt }) + } + + @Test + fun entriesNeedNotArriveSorted() = + runTest { + val shuffled = NegentropyLocalIndex.of(listOf(idAt(1005), idAt(1001), idAt(1009), idAt(1003))) + assertEquals(2, shuffled.count(window(1001, 1003))) + assertEquals(listOf(1001L, 1003L), shuffled.entriesFor(window(1001, 1003)).map { it.createdAt }) + } + + @Test + fun theEmptyIndexAnswersZeroForEveryWindow() = + runTest { + assertEquals(0, NegentropyLocalIndex.Empty.count(window(1000, 2000))) + assertTrue(NegentropyLocalIndex.Empty.entriesFor(window(1000, 2000)).isEmpty()) + assertEquals(0, NegentropyLocalIndex.of(emptyList()).count(window(null, null))) + } +} diff --git a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/NostrClientNegentropySyncTest.kt b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/NostrClientNegentropySyncTest.kt index 52aea9a8b4..6f4b121146 100644 --- a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/NostrClientNegentropySyncTest.kt +++ b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/NostrClientNegentropySyncTest.kt @@ -26,6 +26,7 @@ import com.vitorpamplona.geode.testing.RelayClientTest import com.vitorpamplona.geode.testing.preload import com.vitorpamplona.quartz.nip01Core.core.Event import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient +import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.NegentropyLocalIndex import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.NegentropySyncException import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAllPages import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.negentropySync @@ -36,6 +37,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer import com.vitorpamplona.quartz.nip01Core.relay.server.policies.PassThroughPolicy import com.vitorpamplona.quartz.nip01Core.relay.server.policies.PolicyResult +import com.vitorpamplona.quartz.nip01Core.store.IdAndTime import com.vitorpamplona.quartz.nip77Negentropy.NegentropySettings import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers @@ -48,6 +50,8 @@ import kotlin.test.Test import kotlin.test.assertEquals import kotlin.test.assertFailsWith import kotlin.test.assertFalse +import kotlin.test.assertNotNull +import kotlin.test.assertNull import kotlin.test.assertTrue class NostrClientNegentropySyncTest : RelayClientTest() { @@ -252,6 +256,11 @@ class NostrClientNegentropySyncTest : RelayClientTest() { * The "try negentropy, else page" combinator: against the same over-cap relay * where raw [negentropySync] throws, [negentropySyncOrFetch] transparently pages * and delivers every event, reporting that it fell back. + * + * Every event here shares one `created_at`, so the whole filter IS the + * un-reconcilable window: it is drained as one paged window rather than by + * abandoning the sync, which is why `fallbackCause` is null. On a filter + * spanning more than this second, everything outside it still reconciles. */ @Test fun orFetchPagesWhenNegentropyCannotReconcile() = @@ -275,11 +284,9 @@ class NostrClientNegentropySyncTest : RelayClientTest() { assertEquals(10, got.map { it.id }.toSet().size, "all events delivered via the paging fallback") assertEquals(10, result.downloaded) - assertTrue(result.pagedFallback, "it should have fallen back to paging") - assertEquals( - NegentropySyncException.Reason.OVER_MAX_SYNC_EVENTS, - result.fallbackCause?.reason, - ) + assertTrue(result.pagedFallback, "part of the range came over REQ, so this was not a clean reconcile") + assertEquals(1, result.pagedWindows, "exactly the one un-reconcilable window was paged") + assertNull(result.fallbackCause, "the sync was not abandoned — one window was drained by paging") } finally { client.disconnect() scope.cancel() @@ -375,4 +382,156 @@ class NostrClientNegentropySyncTest : RelayClientTest() { assertFalse(result.pagedFallback, "negentropy should have handled it") assertEquals(8, result.negentropy?.downloaded) } + + /** + * The caller's own count splits a window BEFORE the relay is asked for it. + * + * Nothing here overflows — the relay would have reconciled the whole filter + * in one NEG-OPEN — so every split is driven by [NegentropyLocalIndex.count] + * against `targetWindow`. That is what bounds the entries a caller has to + * materialise: without it the first (and only) window is the whole filter, + * and the local set for it is the whole corpus. + */ + @Test + fun targetWindowSplitsFromTheLocalCountAlone() = + runBlocking { + // 40 seconds of history, one event each; we already hold the even ones. + val all = (0 until 40).map { SyntheticEvents.fakeEvent(idSeed = it + 1, kind = 1, createdAt = 1000L + it) } + defaultRelay.preload(all) + val ours = all.filterIndexed { i, _ -> i % 2 == 0 }.map { IdAndTime(it.createdAt, it.id) } + + val asked = mutableListOf() + val index = + object : NegentropyLocalIndex { + val inner = NegentropyLocalIndex.of(ours) + + override suspend fun count(window: Filter): Int { + asked += window + return inner.count(window) ?: 0 + } + + override suspend fun entriesFor(window: Filter) = inner.entriesFor(window) + } + + val got = mutableListOf() + val result = + withTimeout(60_000) { + client.negentropySync( + relay = defaultRelayUrl, + filter = Filter(kinds = listOf(1)), + localIndex = index, + targetWindow = 5, + ) { got.add(it) } + } + + assertEquals(20, got.map { it.id }.toSet().size, "only the half we lacked comes down") + assertTrue(result.windows > 1, "the local count alone must have split the filter") + assertTrue(asked.isNotEmpty(), "windows must be counted before they are asked for") + assertNull(result.peerCap, "nothing was refused, so there is no cap to report") + } + + /** Passing no target keeps the old shape: one window until the relay objects. */ + @Test + fun withoutATargetTheLocalCountIsNeverConsulted() = + runBlocking { + defaultRelay.preload(SyntheticEvents.batch(20, kind = 1)) + var counted = 0 + val index = + object : NegentropyLocalIndex { + override suspend fun count(window: Filter): Int { + counted++ + return 1_000_000 + } + + override suspend fun entriesFor(window: Filter) = emptyList() + } + + val result = + withTimeout(20_000) { + client.negentropySync( + relay = defaultRelayUrl, + filter = Filter(kinds = listOf(1)), + localIndex = index, + ) { } + } + + assertEquals(0, counted, "targetWindow = 0 must not ask the store anything") + assertEquals(1, result.windows) + assertEquals(20, result.downloaded) + } + + /** + * A relay that refuses for size states its cap, and the client reports it — + * so the next sync can start at a window that fits instead of rediscovering + * it by halving. + */ + @Test + fun theRelaysCapIsReportedBack() = + runBlocking { + val hub = InProcessRelays(negentropySettings = NegentropySettings(maxSyncEvents = 3)) + val scope = CoroutineScope(Dispatchers.Default + SupervisorJob()) + val client = NostrClient(hub, scope) + try { + val url = RelayUrlNormalizer.normalize("ws://127.0.0.1:7786/") + hub.getOrCreate(url).preload((0 until 12).map { SyntheticEvents.fakeEvent(idSeed = it + 1, kind = 1, createdAt = 1000L + it) }) + + val result = + withTimeout(60_000) { + client.negentropySync(relay = url, filter = Filter(kinds = listOf(1))) { } + } + + assertEquals(12, result.downloaded) + assertTrue(result.windows > 1) + assertEquals(3L, result.peerCap, "the relay stated its own max_sync_events") + } finally { + client.disconnect() + scope.cancel() + hub.close() + } + } + + /** + * One second the relay will not reconcile at any window size costs that + * second, not the sync. + * + * The whole point of the [NegentropyOrFetchResult.pagedWindows] path: the + * dense second is drained over REQ while everything around it still + * reconciles. Before, the exception from that one window abandoned the whole + * sync and re-paged the entire filter — on a large corpus, exactly the cost + * negentropy was there to avoid. + */ + @Test + fun oneUnreconcilableSecondDoesNotCostTheRestOfTheFilter() = + runBlocking { + val hub = InProcessRelays(negentropySettings = NegentropySettings(maxSyncEvents = 3)) + val scope = CoroutineScope(Dispatchers.Default + SupervisorJob()) + val client = NostrClient(hub, scope) + try { + val url = RelayUrlNormalizer.normalize("ws://127.0.0.1:7787/") + // Ten events crammed into one second — no created_at window can + // separate them — plus five ordinary seconds around them. + val dense = (1..10).map { SyntheticEvents.fakeEvent(idSeed = it, kind = 1, createdAt = 1000L) } + val sparse = (0 until 5).map { SyntheticEvents.fakeEvent(idSeed = 100 + it, kind = 1, createdAt = 2000L + it) } + hub.getOrCreate(url).preload(dense + sparse) + + val got = mutableListOf() + val result = + withTimeout(60_000) { + client.negentropySyncOrFetch( + relay = url, + filter = Filter(kinds = listOf(1)), + ) { got.add(it) } + } + + assertEquals(15, got.map { it.id }.toSet().size, "everything is delivered, by whichever route") + assertEquals(1, result.pagedWindows, "only the dense second is paged") + assertNull(result.fallbackCause, "the sync itself was never abandoned") + val negentropy = assertNotNull(result.negentropy, "the rest of the range still reconciled") + assertTrue(negentropy.windows > 1) + } finally { + client.disconnect() + scope.cancel() + hub.close() + } + } } From 8555309492a6fc54039cc06d7d08fa3077cd09bb Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 6 Aug 2026 20:12:36 +0000 Subject: [PATCH 3/3] negentropy: audit fixes over the windowing change MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A read-back over the two commits before this, rather than a failure — which is the only way these would have turned up, since every one of them lives on a path that runs when something has already gone wrong. **accept() is no longer single-threaded, and its comment said it was.** "Both phases run sequentially, so no concurrent access" was true right up until a paged window started running on a reconciler coroutine while the sync's own delivery consumer was still calling accept(). An unguarded HashSet between two coroutines can corrupt, and the delivered counter can lose updates. Now behind a Mutex — with onEvent kept INSIDE it, because callers are promised it never runs concurrently with itself and some of them keep unsynchronised state in that callback. pagedWindows becomes an AtomicInt for the same reason. **The kotlinx cap parse could take down the whole frame.** `.jsonPrimitive` throws on an object or array, so a relay putting something structured in the fourth element would have failed the NEG-ERR and lost the reason with it — where before that element existed, anything extra was simply ignored. `as?` restores that. Both mappers are now tested against a structured fourth element as well as a string one. **Int overflow in the split fan-out.** `mine + ceiling - 1` wraps when a window holds close to Int.MAX events, which is reachable on exactly the corpora this targets; done in Long now. **The count-driven split cuts N ways, not two.** The work queue is FIFO, so halving means every internal node's count() runs before the first NEG-OPEN goes out: on a corpus ~30,000 windows wide that is ~30,000 store counts of dead time with nothing downloading. Cutting into ceil(count/budget) pieces (capped at 32) reaches the same corpus in about three levels instead of fifteen, and pieces that guess wrong are re-split by the same rule. **The budget moves by CAS.** With reconcileConcurrency > 1 two reconcilers adjust it at once, and a lost SHRINK is the one that costs something real: the next window is then asked at a size the relay has already refused. --- .../kotlinSerialization/MessageKSerializer.kt | 6 +- .../client/accessories/NegentropyStoreSync.kt | 47 +++++- .../NostrClientNegentropySyncExt.kt | 143 +++++++++++++----- .../nip77Negentropy/Nip77SerializationTest.kt | 18 +++ 4 files changed, 173 insertions(+), 41 deletions(-) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/kotlinSerialization/MessageKSerializer.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/kotlinSerialization/MessageKSerializer.kt index 93c0c360cb..c0a7972e10 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/kotlinSerialization/MessageKSerializer.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/kotlinSerialization/MessageKSerializer.kt @@ -191,8 +191,10 @@ object MessageKSerializer : KSerializer { reason = if (array.size > 2) array[2].jsonPrimitive.content else "", // Optional, and only a number: a relay that puts something // else there is telling us nothing rather than breaking the - // frame. - cap = if (array.size > 3) array[3].jsonPrimitive.longOrNull else null, + // frame. `as?` rather than `.jsonPrimitive`, which THROWS on + // an object or array — that would fail the whole message and + // lose the reason, where before this element was ignored. + cap = if (array.size > 3) (array[3] as? JsonPrimitive)?.longOrNull else null, ) } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NegentropyStoreSync.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NegentropyStoreSync.kt index a0272abecc..02bb576509 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NegentropyStoreSync.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NegentropyStoreSync.kt @@ -26,6 +26,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl import com.vitorpamplona.quartz.nip01Core.store.IEventStore +import com.vitorpamplona.quartz.nip01Core.store.IdAndTime import com.vitorpamplona.quartz.nip01Core.store.verifyAndInsert import kotlinx.coroutines.async import kotlinx.coroutines.awaitAll @@ -92,6 +93,15 @@ class NegentropyStoreSync( * @param concurrency relays synced at once by [sync] (a relay's own filters stay sequential). * @param idleTimeoutMs idle watchdog for reconciles / fetches / pages. * @param publishTimeoutSecs OK-confirmation wait per uploaded event. + * @param targetWindow events per reconcile window, or `0` to snapshot the + * whole filter up front (the default, and what this class always did). + * + * Above zero, the store is read one `created_at` window at a time through + * a [NegentropyLocalIndex] instead: the id snapshot stops being O(matched + * set) — it is the largest thing this class holds — at the price of an + * indexed count + range read per window. Worth turning on exactly when the + * filter matches more than fits comfortably in memory; pointless below + * that, where one snapshot shared by the whole group is cheaper. */ class Config( val down: Boolean = true, @@ -105,6 +115,7 @@ class NegentropyStoreSync( val concurrency: Int = 4, val idleTimeoutMs: Long = 30_000L, val publishTimeoutSecs: Long = 15, + val targetWindow: Int = 0, ) /** Outcome of one `(relay, filter)` group. `error` is null on success. */ @@ -166,7 +177,14 @@ class NegentropyStoreSync( // events (~40 B/entry vs ~1 KB), which matters when a relay hosts a large // matched set. The events the reconcile decides to UP-publish (the small // residual haves) are fetched by id on demand in the uploader below. - val localEntries = store.snapshotIdsForNegentropy(listOf(filter)) + // + // With a targetWindow, even those 40 B/entry are read per window rather + // than for the whole filter — on a large store that snapshot is the + // biggest thing this class allocates, and it is allocated before the + // first frame goes out. + val windowed = config.targetWindow > 0 + val localIndex = if (windowed) StoreWindowIndex(store) else null + val localEntries = if (windowed) emptyList() else store.snapshotIdsForNegentropy(listOf(filter)) val downloaded = AtomicInt(0) val uploaded = AtomicInt(0) @@ -210,6 +228,8 @@ class NegentropyStoreSync( relay = relay, filter = filter, localEntries = localEntries, + localIndex = localIndex, + targetWindow = config.targetWindow, batchSize = config.idChunk, idleTimeoutMs = config.idleTimeoutMs, reconcileConcurrency = config.reconcileConcurrency, @@ -311,3 +331,28 @@ class NegentropyStoreSync( return stored.load() } } + +/** + * [NegentropyLocalIndex] over an [IEventStore]: the window engine's per-window + * reads answered straight from the store's `created_at` index. + * + * The windows handed here are the caller's own filter with `since`/`until` + * narrowed, so they can go to the store as-is. A count the store cannot answer + * comes back null rather than throwing — the engine then simply stops + * pre-splitting that window and lets the relay's refusal decide, which is the + * behaviour without an index at all. + */ +private class StoreWindowIndex( + private val store: IEventStore, +) : NegentropyLocalIndex { + override suspend fun count(window: Filter): Int? = + try { + store.count(window) + } catch (e: CancellationException) { + throw e + } catch (_: Exception) { + null + } + + override suspend fun entriesFor(window: Filter): List = store.snapshotIdsForNegentropy(listOf(window)) +} diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncExt.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncExt.kt index 6595cc763e..51941e1856 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncExt.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncExt.kt @@ -329,6 +329,7 @@ class NegentropyOrFetchResult( * Use [negentropySync] directly if you want to decide the fallback yourself (try * another relay, narrow the filter, abort, …) instead of always paging. */ +@OptIn(ExperimentalAtomicApi::class) suspend fun INostrClient.negentropySyncOrFetch( relay: NormalizedRelayUrl, filter: Filter, @@ -346,18 +347,29 @@ suspend fun INostrClient.negentropySyncOrFetch( ): NegentropyOrFetchResult { val seen = HashSet() var delivered = 0 - var pagedWindows = 0 + val pagedWindows = AtomicInt(0) - // Shared dedup + cap across both phases. Returns true if the event was new and - // delivered. Both phases run sequentially, so no concurrent access. - suspend fun accept(event: Event): Boolean { - if ((maxEvents <= 0 || delivered < maxEvents) && seen.add(event.id)) { - delivered++ - onEvent(event) - return true + // Shared dedup + cap across every path that delivers. + // + // The lock is not optional. The two phases used to run strictly one after + // the other, but a paged window now runs DURING the negentropy phase, on a + // reconciler coroutine, while the sync's own delivery consumer is calling + // this too — an unguarded HashSet between them can corrupt, and the count + // can lose updates. onEvent stays INSIDE the lock deliberately: callers are + // promised it never runs concurrently with itself, and some of them keep + // unsynchronised state in it. + val gate = Mutex() + + suspend fun accept(event: Event): Boolean = + gate.withLock { + if ((maxEvents <= 0 || delivered < maxEvents) && seen.add(event.id)) { + delivered++ + onEvent(event) + true + } else { + false + } } - return false - } return try { val result = @@ -379,7 +391,7 @@ suspend fun INostrClient.negentropySyncOrFetch( // already reconciled cleanly walked again over REQ, which on a // large corpus is the entire cost negentropy was there to save. onUnreconcilableWindow = { window -> - pagedWindows++ + pagedWindows.incrementAndFetch() val pageTimeoutMs = if (idleTimeoutMs > 0) idleTimeoutMs else DEFAULT_DOWNLOAD_IDLE_MS fetchAllPages(relay, listOf(window), pageTimeoutMs) { event -> if (accept(event)) onProgress?.invoke(delivered, delivered) @@ -392,10 +404,10 @@ suspend fun INostrClient.negentropySyncOrFetch( // Any paged window makes this not a clean reconcile — see the // property doc: under-reporting it would let a caller record // coverage it never compared. - pagedFallback = pagedWindows > 0, + pagedFallback = pagedWindows.load() > 0, negentropy = result, fallbackCause = null, - pagedWindows = pagedWindows, + pagedWindows = pagedWindows.load(), ) } catch (e: NegentropySyncException) { // Negentropy couldn't enumerate the set — page the whole filter instead, @@ -411,7 +423,7 @@ suspend fun INostrClient.negentropySyncOrFetch( pagedFallback = true, negentropy = null, fallbackCause = e, - pagedWindows = pagedWindows, + pagedWindows = pagedWindows.load(), ) } } @@ -567,7 +579,9 @@ internal suspend fun reconcileWindows( onPeerCap: ((Long) -> Unit)? = null, // Given a minimal window the relay will not reconcile at any size, instead // of throwing. The caller drains it however it can (paging it over REQ) and - // the sweep carries on with the rest of the filter. + // the sweep carries on with the rest of the filter. It runs ON the reconciler + // that hit the window, so a slow drain holds that reconciler — with + // reconcileConcurrency = 1 the rest of the sweep waits for it. onUnreconcilableWindow: (suspend (Filter) -> Unit)? = null, sendNeedBatch: suspend (List) -> Unit, sendHaveBatch: (suspend (List) -> Unit)?, @@ -597,25 +611,57 @@ internal suspend fun reconcileWindows( // targetWindow at 0 nothing below this line does anything. val budget = AtomicInt(targetWindow) - // Splits a window in two and queues both halves. Returns false when the - // window is already minimal — `created_at` is in seconds, so that is the - // floor, not a tuning choice. - fun splitInto( + // Every budget move goes through here. With reconcileConcurrency > 1 two + // reconcilers adjust it at once, and read-then-store can drop one of them — + // a lost SHRINK being the one that costs something real, since the next + // window is then asked at a size the relay has already refused. + fun budgetTo(next: (Int) -> Int) { + while (true) { + val now = budget.load() + val want = next(now) + if (want == now || budget.compareAndSet(now, want)) return + } + } + + /** + * Cuts `[lo, hi]` into [pieces] equal spans of time and queues them all. A + * window already at the floor is left alone — `created_at` is in seconds, so + * that is where splitting ends, not a tuning choice. Both callers check that + * themselves; the guard here is so a third one cannot silently lose a window. + * + * [pieces] > 2 exists for the count-driven split, where we know HOW FAR over + * the budget a window is and can land near the right size in one step. + * Halving instead costs a store count per level of a tree that can be ~15 + * deep on a large corpus, and — since the queue is FIFO — every one of those + * counts happens before the first window is reconciled at all. + */ + suspend fun splitInto( pendingWindow: Filter, lo: Long, hi: Long, - ): Boolean { - if (hi - lo <= MIN_WINDOW_SECONDS) return false - val mid = lo + (hi - lo) / 2 - remaining.incrementAndFetch() - // The lower child gets the finite midpoint; the upper child KEEPS this - // window's original `until` (which may be null = unbounded). Replacing - // null with `now()` here would drop every event dated after now() - // (clock skew) once any split happens, while the un-split path would - // have included them. - pending.trySend(pendingWindow.copy(since = lo, until = mid)) - pending.trySend(pendingWindow.copy(since = mid + 1, until = pendingWindow.until)) - return true + pieces: Int = 2, + ) { + if (hi - lo <= MIN_WINDOW_SECONDS) return + val span = hi - lo + 1 + // Never more pieces than there are seconds to give them. + val n = pieces.toLong().coerceIn(2L, minOf(span, MAX_SPLIT_FANOUT.toLong())).toInt() + val step = span / n + remaining.addAndFetch(n - 1) + var start = lo + repeat(n) { i -> + val last = i == n - 1 + // The top piece KEEPS this window's original `until` (which may be + // null = unbounded). Replacing null with `now()` here would drop + // every event dated after now() (clock skew) once any split happens, + // while the un-split path would have included them. + if (last) { + pending.send(pendingWindow.copy(since = start, until = pendingWindow.until)) + } else { + val end = start + step - 1 + pending.send(pendingWindow.copy(since = start, until = end)) + start = end + 1 + } + } } val reconcilers = @@ -636,7 +682,14 @@ internal suspend fun reconcileWindows( if (ceiling > 0 && hi - lo > MIN_WINDOW_SECONDS) { val mine = local.count(window) if (mine != null && mine > ceiling) { - splitInto(window, lo, hi) + // How many windows this one is worth, not just "two": + // the count says how far over budget we are, and + // uneven density is corrected by the same check on + // each piece. + // Long arithmetic: `mine` can be near Int.MAX on a + // corpus this size, and the +ceiling would wrap. + val over = (mine.toLong() + ceiling - 1) / ceiling + splitInto(window, lo, hi, pieces = over.coerceAtMost(MAX_SPLIT_FANOUT.toLong()).toInt()) continue } } @@ -662,9 +715,12 @@ internal suspend fun reconcileWindows( // asked for, so a sync that met one dense stretch // does not stay small for the rest of the timeline. if (targetWindow > 0) { - val now = budget.load() - if (now < targetWindow) { - budget.store(minOf(targetWindow, (now * BUDGET_GROWTH).toInt().coerceAtLeast(now + 1))) + budgetTo { now -> + if (now >= targetWindow) { + now + } else { + minOf(targetWindow, (now * BUDGET_GROWTH).toInt().coerceAtLeast(now + 1)) + } } } if (remaining.decrementAndFetch() == 0) pending.close() @@ -677,13 +733,16 @@ internal suspend fun reconcileWindows( outcome.cap?.let { cap -> onPeerCap?.invoke(cap) if (targetWindow > 0) { - val fitted = (cap * CAP_MARGIN).toInt().coerceAtLeast(1) - if (fitted < budget.load()) budget.store(fitted) + val fitted = + (cap * CAP_MARGIN) + .coerceIn(1.0, Int.MAX_VALUE.toDouble()) + .toInt() + budgetTo { now -> minOf(now, fitted) } } } if (outcome.cap == null && targetWindow > 0) { // No number to go on: halve and find out. - budget.store((budget.load() / 2).coerceAtLeast(1)) + budgetTo { now -> (now / 2).coerceAtLeast(1) } } if (hi - lo <= MIN_WINDOW_SECONDS) { // A minimal window that still overflows: @@ -1295,6 +1354,14 @@ private const val MIN_WINDOW_SECONDS = 1L */ private const val MAX_WINDOWS = 100_000 +/** + * Most pieces one count-driven split may cut a window into. Bounds both the + * queue and the depth: with 32, a corpus 30,000 windows wide is reached in + * three levels instead of fifteen, and the pieces that guessed wrong are + * re-split by the same rule. + */ +private const val MAX_SPLIT_FANOUT = 32 + /** * How much of a relay's stated `max_sync_events` a window actually aims for. * The margin absorbs what the relay gains between stating that number and diff --git a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip77Negentropy/Nip77SerializationTest.kt b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip77Negentropy/Nip77SerializationTest.kt index 7a71cc9a1b..d8d0bcc50b 100644 --- a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip77Negentropy/Nip77SerializationTest.kt +++ b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip77Negentropy/Nip77SerializationTest.kt @@ -175,6 +175,24 @@ class Nip77SerializationTest { assertEquals(null, kotlin.cap) } + @Test + fun deserializeNegErrMessageWithStructuredFourthElement_bothMappers() { + // A fourth element that is an object or array must degrade to no cap, + // NOT fail the frame — the reason is the part that matters, and before + // this element existed any extra was simply ignored. + val json = """["NEG-ERR","neg-sub1","blocked: too many query results",{"max":10}]""" + + val jackson = JacksonMapper.fromJsonToMessage(json) + assertTrue(jackson is NegErrMessage) + assertEquals("blocked: too many query results", jackson.reason) + assertEquals(null, jackson.cap) + + val kotlin = KotlinSerializationMapper.fromJsonToMessage(json) + assertTrue(kotlin is NegErrMessage) + assertEquals("blocked: too many query results", kotlin.reason) + assertEquals(null, kotlin.cap) + } + @Test fun negErrMessageWithCap_crossDeserialization() { val msg = NegErrMessage("neg-sub1", "blocked: too many records", 500_000L)