perf(marmot): stop the KeyPackage lookup shortly after the first package

Adding a member still took most of a minute after the commit fix. The
KeyPackage lookup drains every relay to EOSE to pick the newest package,
and it asks the invitee's relays plus ours, so one relay that never sends
EOSE held each invite for the full 30s idle window.

Once a KeyPackage arrives, the other relays get 3s to report a newer one
and the fetch stops, returning the newest seen. A rotation visible only
on a slower relay can be missed; the older package still opens.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
Vitor Pamplona
2026-10-05 15:19:18 -04:00
co-authored by Claude Opus 5.5
parent 98a32f4b68
commit 930466e9b0
2 changed files with 158 additions and 5 deletions
@@ -23,8 +23,12 @@ package com.vitorpamplona.quartz.marmot.mip00KeyPackages
import com.vitorpamplona.quartz.marmot.MarmotFilters
import com.vitorpamplona.quartz.nip01Core.core.HexKey
import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAll
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAllWithHooks
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import kotlinx.coroutines.CompletableDeferred
import kotlinx.coroutines.coroutineScope
import kotlinx.coroutines.delay
import kotlinx.coroutines.launch
/**
* Discovery helpers for MIP-00 KeyPackages.
@@ -73,19 +77,49 @@ object KeyPackageFetcher {
* relays and pick whichever replied first, which would frequently be the
* older event. Draining to EOSE and selecting by `created_at` matches
* MDK/whitenoise semantics and keeps freshly-rotated bundles reachable.
*
* Draining to EOSE is bounded by [settleAfterFirstMs], though. The relay
* set unions the invitee's relays with ours, and one of them that never
* sends EOSE held every invite for the whole [idleTimeoutMs] — adding a
* member took most of a minute. Once a KeyPackage has arrived, the other
* relays get [settleAfterFirstMs] to report a newer one and the fetch
* stops. A relay slower than that can only cost us a rotation that
* happened in the last moments, and the older package still opens.
*/
suspend fun fetchKeyPackage(
client: INostrClient,
targetPubKey: HexKey,
relays: Set<NormalizedRelayUrl>,
idleTimeoutMs: Long = 30_000,
settleAfterFirstMs: Long = 3_000,
): KeyPackageEvent? {
if (relays.isEmpty()) return null
val filter = MarmotFilters.keyPackagesByAuthor(targetPubKey)
val events = client.fetchAll(filters = relays.associateWith { listOf(filter) }, idleTimeoutMs = idleTimeoutMs)
// fetchAll returns events sorted by created_at DESC, so the first
// KeyPackageEvent is the most recent one any relay had.
return events.firstNotNullOfOrNull { it as? KeyPackageEvent }
// Collected from inside onEvent (single-threaded) rather than read from the
// return value, which a cancelled fetch discards.
val found = mutableListOf<KeyPackageEvent>()
val firstArrived = CompletableDeferred<Unit>()
coroutineScope {
val fetch =
launch {
client.fetchAllWithHooks(filters = relays.associateWith { listOf(filter) }, idleTimeoutMs = idleTimeoutMs) { _, event ->
if (event is KeyPackageEvent) {
found.add(event)
firstArrived.complete(Unit)
}
true
}
}
val settle =
launch {
firstArrived.await()
delay(settleAfterFirstMs)
fetch.cancel()
}
fetch.join()
settle.cancel()
}
return found.maxByOrNull { it.createdAt }
}
/**
@@ -0,0 +1,119 @@
/*
* 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.marmot
import com.vitorpamplona.geode.InProcessRelays
import com.vitorpamplona.quartz.marmot.mip00KeyPackages.KeyPackageEvent
import com.vitorpamplona.quartz.marmot.mip00KeyPackages.KeyPackageFetcher
import com.vitorpamplona.quartz.nip01Core.crypto.KeyPair
import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.publishAndConfirm
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebSocket
import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebSocketListener
import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebsocketBuilder
import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerInternal
import com.vitorpamplona.quartz.utils.TimeUtils
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel
import kotlinx.coroutines.runBlocking
import kotlin.test.AfterTest
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertTrue
import kotlin.time.measureTimedValue
/**
* Adding a Marmot member took most of a minute: the KeyPackage lookup drains every relay to EOSE,
* and the set it asks unions the invitee's relays with ours, so one relay that never answers held
* every invite for the whole 30s idle window.
*/
class KeyPackageFetcherSilentRelayTest {
private val hub = InProcessRelays()
private val scope = CoroutineScope(Dispatchers.Default + SupervisorJob())
private val silent: NormalizedRelayUrl = RelayUrlNormalizer.normalize("ws://127.0.0.1:7772/")
@AfterTest
fun tearDown() {
scope.cancel()
hub.close()
}
/** The hub, except [silent] swallows every REQ: no events, no EOSE, no CLOSED. */
private inner class OneSilentRelay : WebsocketBuilder {
override fun build(
url: NormalizedRelayUrl,
out: WebSocketListener,
): WebSocket {
val inner = hub.build(url, out)
if (url != silent) return inner
return object : WebSocket by inner {
override fun send(msg: String): Boolean = if (msg.startsWith("[\"REQ\"")) true else inner.send(msg)
}
}
}
private fun keyPackage(
signer: NostrSignerInternal,
slot: String,
createdAt: Long,
) = runBlocking {
signer.sign(
KeyPackageEvent.build(
keyPackageBase64 = "AA==",
dTagSlot = slot,
keyPackageRef = "00".repeat(32),
relays = listOf(InProcessRelays.DEFAULT_URL),
createdAt = createdAt,
),
)
}
@Test
fun aSilentRelayDoesNotHoldTheLookupAndTheNewestPackageStillWins() =
runBlocking {
val client = NostrClient(OneSilentRelay(), scope)
val invitee = NostrSignerInternal(KeyPair())
val now = TimeUtils.now()
val older = keyPackage(invitee, "a", now - 3600)
val newer = keyPackage(invitee, "b", now)
assertTrue(client.publishAndConfirm(older, setOf(InProcessRelays.DEFAULT_URL)))
assertTrue(client.publishAndConfirm(newer, setOf(InProcessRelays.DEFAULT_URL)))
val (found, elapsed) =
measureTimedValue {
KeyPackageFetcher.fetchKeyPackage(
client = client,
targetPubKey = invitee.pubKey,
relays = setOf(InProcessRelays.DEFAULT_URL, silent),
idleTimeoutMs = 10_000,
settleAfterFirstMs = 500,
)
}
assertEquals(newer.id, found?.id, "the newest KeyPackage is still the one returned")
assertTrue(elapsed.inWholeMilliseconds < 3_000, "stopped shortly after the first package, not at the 10s idle window (took $elapsed)")
client.disconnect()
}
}