From 58f6a0af6b527b5b39984f07776e3ebc6efc823d Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 7 May 2026 01:03:06 +0000 Subject: [PATCH] fix(relay): forward ephemeral events to live subscribers (NIP-01) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit LiveEventStore.query had a race window between emitting EOSE and registering as a SharedFlow collector — any event emitted in that window was lost because newEventStream has replay=0. The race was usually masked by SQLite write latency for persisted kinds, but it fired reliably for ephemeral kinds (20000-29999) where store.insert is a no-op. Per NIP-01 ephemeral events are not persisted but MUST still reach matching active subscriptions. Fixed by registering the live collector before signalling EOSE via Flow.onSubscription. Adds two tests: one proves an ephemeral event reaches an active subscriber, the other proves it isn't persisted (a follow-up REQ returns zero events). --- .../quartz/relay/Nip01ComplianceTest.kt | 58 +++++++++++++++++++ .../nip01Core/relay/server/LiveEventStore.kt | 28 +++++---- 2 files changed, 75 insertions(+), 11 deletions(-) diff --git a/quartz-relay/src/test/kotlin/com/vitorpamplona/quartz/relay/Nip01ComplianceTest.kt b/quartz-relay/src/test/kotlin/com/vitorpamplona/quartz/relay/Nip01ComplianceTest.kt index 3cf9768aec..4a656781ce 100644 --- a/quartz-relay/src/test/kotlin/com/vitorpamplona/quartz/relay/Nip01ComplianceTest.kt +++ b/quartz-relay/src/test/kotlin/com/vitorpamplona/quartz/relay/Nip01ComplianceTest.kt @@ -311,6 +311,64 @@ class Nip01ComplianceTest { client.unsubscribe("live-2") } + // -- Ephemeral events (NIP-01: kinds 20000-29999) ---------------------- + + /** + * Ephemeral events MUST be forwarded to active subscriptions whose + * filters match, even though the relay does not persist them. + */ + @Test + fun ephemeralEventForwardedToActiveSubscription() = + runBlocking { + val ch = Channel(UNLIMITED) + val gotEose = Channel(UNLIMITED) + client.subscribe( + "eph-1", + mapOf(relayUrl to listOf(Filter(kinds = listOf(20_001)))), + object : SubscriptionListener { + override fun onEvent( + event: Event, + isLive: Boolean, + relay: NormalizedRelayUrl, + forFilters: List?, + ) { + ch.trySend(event) + } + + override fun onEose( + relay: NormalizedRelayUrl, + forFilters: List?, + ) { + gotEose.trySend(Unit) + } + }, + ) + withTimeout(5000) { gotEose.receive() } + + hub.getOrCreate(relayUrl).publish(fakeEvent(70, kind = 20_001, content = "ephemeral-payload")) + + val received = withTimeout(5000) { ch.receive() } + assertEquals("ephemeral-payload", received.content) + assertEquals(20_001, received.kind) + client.unsubscribe("eph-1") + } + + /** + * Ephemeral events MUST NOT be persisted. A REQ issued after the + * event was published returns nothing. + */ + @Test + fun ephemeralEventIsNotStoredAndDoesNotShowOnFollowupReq() = + runBlocking { + // Publish ephemeral first — no live subscriber listening. + hub.getOrCreate(relayUrl).publish(fakeEvent(71, kind = 20_002, content = "vanish")) + + // Late subscriber: should see EOSE with no events. + val (events, eose) = collectUntilEose(Filter(kinds = listOf(20_002))) + assertTrue(eose, "EOSE must fire for an ephemeral kind even if zero events match") + assertEquals(0, events.size, "Ephemeral events must not be persisted") + } + // -- Multi-relay -------------------------------------------------------- /** A single client can hold subscriptions against multiple relays simultaneously. */ diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/LiveEventStore.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/LiveEventStore.kt index 5e1cf99996..2b2d87eb73 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/LiveEventStore.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/LiveEventStore.kt @@ -25,6 +25,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter import com.vitorpamplona.quartz.nip01Core.store.IEventStore import kotlinx.coroutines.channels.BufferOverflow import kotlinx.coroutines.flow.MutableSharedFlow +import kotlinx.coroutines.flow.onSubscription /** * A reactive event store that combines historical data retrieval with live event streaming. @@ -56,18 +57,23 @@ class LiveEventStore( onEach: (Event) -> Unit, onEose: () -> Unit, ) { - // 1. Replay stored events matching filters. - store.query(filters, onEach) - - // 2. Signal end of stored events. - onEose() - - // 3. Stream live events until cancelled. - newEventStream.collect { newEvent -> - if (filters.any { it.match(newEvent) }) { - onEach(newEvent) + // Order matters: register the live collector BEFORE replaying + // stored events and signalling EOSE. Otherwise an event emitted + // between EOSE and `collect` is lost because [newEventStream] has + // replay=0. The race is only occasionally visible for kinds the + // store persists (insert latency masks it) but fires reliably for + // ephemeral kinds (20000-29999) where insert is a no-op — and + // ephemeral events MUST still reach matching live subscribers per + // NIP-01. + newEventStream + .onSubscription { + store.query(filters, onEach) + onEose() + }.collect { newEvent -> + if (filters.any { it.match(newEvent) }) { + onEach(newEvent) + } } - } } suspend fun count(filters: List) = store.count(filters)