Suspend the incoming-message chain down to SubscriptionListener.onEvent

A consumer that cannot suspend has to block, and blocking here deadlocks
the whole client.

Measured on a mirror built against this library, twice, ~13 minutes after
each start: all 64 shared coroutine workers parked in `runBlocking` beneath
`trySendBlocking`, called from the websocket message callback. The consumer
draining that channel needed threads from the same pool to reach its store,
so it could never make room, so the producers never woke. Every stream, the
health reporter, all of it stopped, at 2% CPU with a healthy, idle backend.
A full queue was the symptom; producers eating the threads the drain needed
was the cause.

The coroutine context was already there — BasicOkHttpWebSocket has always
processed messages inside `scope.launch { for (message in incomingMessages) }`
— so the only thing forcing a blocking hand-off was that the hops in between
were declared non-suspend. Now they are not:

    WebSocketListener.onMessage
    RelayConnectionListener.onIncomingMessage
    PoolRequests/PoolCounts/PoolEventOutbox.onIncomingMessage
    SubscriptionListener.onEvent
    fetchAllPages / negentropy accessories' onEvent parameter

A consumer that fills its buffer now suspends and releases its thread rather
than holding it, which is the same reasoning BasicOkHttpWebSocket already
documents for keeping its own channel UNLIMITED so a slow consumer cannot
block OkHttp reader threads. This extends it one layer down.

BLE is the one transport whose callback genuinely cannot suspend — the
platform hands notifications to a plain callback — so BleNostrClient gets
the same treatment the websocket transport already had: an UNLIMITED
hand-off channel so the BLE stack is never blocked, drained by ONE coroutine
so message order survives the boundary.

Tests that drove these entry points directly now do so from `runTest`, or
from `runBlocking` where the call sits inside a raw thread or Runnable that
models a platform callback.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Vitor Pamplona
2026-08-05 12:45:02 -04:00
co-authored by Claude Opus 5
parent e822911d09
commit a3fac4fc08
124 changed files with 770 additions and 682 deletions
@@ -80,7 +80,7 @@ class NappletLiveSubscriptions {
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -113,7 +113,7 @@ object ClinkDebitPayer {
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -85,7 +85,7 @@ object ClinkOfferPayer {
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -158,7 +158,7 @@ class BootRelayDiagnostics(
}
}
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -112,7 +112,7 @@ class DmRelayDiagnosticsLogger(
Log.d(TAG) { "[+${at()}ms] REQ -> ${relay.url.url} success=$success ${cmdStr.take(400)}" }
}
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -78,7 +78,7 @@ abstract class PerUniqueIdEoseManager<T, U : Any>(
newEose(key, relay, TimeUtils.now(), forFilters)
}
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -90,7 +90,7 @@ abstract class PerUserAndFollowListEoseManager<T, U : Any>(
newEose(key, relay, TimeUtils.now(), forFilters)
}
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -77,7 +77,7 @@ abstract class PerUserEoseManager<T>(
newEose(key, relay, TimeUtils.now(), forFilters)
}
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -53,7 +53,7 @@ abstract class SingleSubNoEoseCacheEoseManager<T>(
}
}
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -78,7 +78,7 @@ class NotifyCoordinator(
}
}
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -116,7 +116,7 @@ class AccountFollowsLoaderSubAssembler(
newEose(TimeUtils.now(), relay, forFilters)
}
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -193,7 +193,7 @@ class AccountNotificationsHistoryEoseManager(
// cursors so a late callback can't move another account's cursors. newEose runs regardless.
val myCursors = key.account.notificationHistory
return object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -124,7 +124,7 @@ class NwcNotificationsEoseManager(
newEose(key, relay, TimeUtils.now(), forFilters)
}
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -115,7 +115,7 @@ class AccountGiftWrapsHistoryEoseManager(
// cursors so a late callback can't move another account's cursors. newEose runs regardless.
val myCursors = key.account.chatroomList.giftWrapHistory
return object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -74,7 +74,7 @@ class UserWatcherSubAssembler(
newEose(relay, TimeUtils.now(), forFilters)
}
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -42,7 +42,7 @@ class RelaySpeedLogger(
private val clientListener =
object : RelayConnectionListener {
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -48,7 +48,7 @@ class RelayUsageListener(
}
}
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -116,7 +116,7 @@ class ChatroomNip04HistorySubAssembler(
// so a late callback can't move another room's cursors. newEose (framework bookkeeping) runs anyway.
val myCursors = cursorsFor(key)
return object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -170,7 +170,7 @@ class ConcordChannelHistorySubAssembler(
// cursors so a late callback can't move another channel's cursors. newEose runs regardless.
val myCursors = cursorsFor(key)
return object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -126,7 +126,7 @@ class RelayGroupOpenChatHistorySubAssembler(
// cursors so a late callback can't move another group's cursors. newEose runs regardless.
val myCursors = cursorsFor(key)
return object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -123,7 +123,7 @@ class RelayGroupOpenThreadsHistorySubAssembler(
// cursors so a late callback can't move another group's cursors. newEose runs regardless.
val myCursors = cursorsFor(key)
return object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -108,7 +108,7 @@ class ChatroomListNip04HistorySubAssembler(
// cursors so a late callback can't move another account's cursors. newEose runs regardless.
val myCursors = key.account.chatroomList.nip04History
return object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -70,7 +70,7 @@ class ChessFeedFilterSubAssembler(
newEose(key, relay, TimeUtils.now(), forFilters)
}
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -456,7 +456,7 @@ class EventSync(
}
}
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -668,7 +668,7 @@ class Context(
val filters = relays.associateWith { listOf(responseFilter) }
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -133,7 +133,7 @@ object GeochatCommands {
val subId = newSubId()
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -167,7 +167,7 @@ object NipCommand {
val remaining = SEARCH_RELAYS.toMutableSet()
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -118,7 +118,7 @@ object NostrConnect {
val subId = newSubId()
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -81,7 +81,7 @@ object SubscribeCommand {
val subId = newSubId()
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -100,7 +100,7 @@ class ChessRelayFetchHelper(
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -99,7 +99,7 @@ class FeedMetadataCoordinator(
val listener =
if (onEvent != null) {
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -295,7 +295,7 @@ class FeedMetadataCoordinator(
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -371,7 +371,7 @@ class FeedMetadataCoordinator(
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -90,7 +90,7 @@ abstract class PerKeyEoseManager<T, K : Any>(
newEose(queryState, relay, TimeUtils.now(), forFilters)
}
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -86,7 +86,7 @@ abstract class SingleSubEoseManager<T>(
newEose(relay, TimeUtils.now(), forFilters)
}
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -44,7 +44,7 @@ class RelayHealthListener(
store.recordConnect(relay.url, TimeUtils.now())
}
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -118,7 +118,7 @@ class BroadcastTracker {
}
}
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -294,7 +294,7 @@ class BroadcastTracker {
}
}
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -371,7 +371,7 @@ class OutboxDispatcher(
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -415,7 +415,7 @@ class OutboxDispatcher(
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -90,7 +90,7 @@ class ChessEventBroadcaster(
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -212,7 +212,7 @@ fun WindowLoadTracker.trackingListener(forward: (NormalizedRelayUrl, List<Filter
forward(relay, forFilters)
}
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -53,7 +53,7 @@ class RelayLatencyListener(
tracker.recordSent(relay.url, cmd, success)
}
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -97,7 +97,7 @@ class DmInboxRelayResolverOutboxTest {
filterList.forEach { filter ->
filter.kinds?.forEach { kind ->
script[kind to relay]?.forEach { event ->
listener?.onEvent(event, isLive = false, relay = relay, forFilters = null)
kotlinx.coroutines.runBlocking { listener?.onEvent(event, isLive = false, relay = relay, forFilters = null) }
}
}
}
@@ -171,7 +171,7 @@ class OutboxDispatcherTest {
filterList.forEach { filter ->
filter.kinds?.forEach { kind ->
script[kind to relay]?.forEach { event ->
listener?.onEvent(event, isLive = false, relay = relay, forFilters = null)
kotlinx.coroutines.runBlocking { listener?.onEvent(event, isLive = false, relay = relay, forFilters = null) }
}
}
}
@@ -1786,7 +1786,7 @@ fun MainContent(
filters = listOf(filter),
listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: com.vitorpamplona.quartz.nip01Core.core.Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -1845,7 +1845,7 @@ fun MainContent(
relays = outbox,
listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: com.vitorpamplona.quartz.nip01Core.core.Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -249,7 +249,7 @@ class AccountManager internal constructor(
}
}
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -164,7 +164,7 @@ class FollowPacksState(
private fun subscribeToDiscovery() {
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: com.vitorpamplona.quartz.nip01Core.core.Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -45,7 +45,7 @@ fun RelayConnectionManager.subscribeMetadataFor(
val filter = Filter(kinds = listOf(MetadataEvent.KIND), authors = pubkeys)
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -92,7 +92,7 @@ fun FromThePackFeed(
listOf(filter),
listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -102,7 +102,7 @@ fun RenderFollowPackCard(
listOf(filter),
listener =
object : com.vitorpamplona.quartz.nip01Core.relay.client.reqs.SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: com.vitorpamplona.quartz.nip01Core.core.Event,
isLive: Boolean,
relay: com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl,
@@ -200,7 +200,7 @@ open class RelayConnectionManager(
filters = filterMap,
listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -264,7 +264,7 @@ open class RelayConnectionManager(
}
}
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -69,7 +69,7 @@ class DesktopRelayUserSearchDelegate(
relays = relays,
listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -96,7 +96,7 @@ class DesktopChessSubscriptionController(
relays = state.relays,
listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -303,7 +303,7 @@ class DesktopRelaySubscriptionsCoordinator(
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -428,7 +428,7 @@ class DesktopRelaySubscriptionsCoordinator(
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -81,7 +81,7 @@ fun rememberSubscription(
relays = cfg.relays,
listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -225,7 +225,7 @@ fun ImportFollowListDialog(
relays = relays,
listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -277,7 +277,7 @@ fun ImportFollowListDialog(
),
listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -858,7 +858,7 @@ private suspend fun fetchMetadataForUsers(
relays = relays,
listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -1666,7 +1666,7 @@ private suspend fun fetchUserLightningAddress(
relays = relays,
listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -175,7 +175,7 @@ class ChatroomListState(
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -139,7 +139,7 @@ object LaunchScenario {
),
listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: com.vitorpamplona.quartz.nip01Core.core.Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -73,7 +73,7 @@ class LaunchFixtureRelayTest {
),
listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -72,7 +72,7 @@ class SubscribeBeforeConnectTest {
),
listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -686,7 +686,7 @@ class MirrorWorker(
val watermark = AtomicLong(initialSince)
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -125,7 +125,7 @@ class GracefulShutdownTest {
val gotEose = Channel<Unit>(UNLIMITED)
val listener =
object : RelayConnectionListener {
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -389,7 +389,7 @@ class KtorRelayTest {
"close-test",
mapOf(server.url.normalizeRelayUrl() to listOf(Filter(kinds = listOf(1)))),
object : com.vitorpamplona.quartz.nip01Core.relay.client.reqs.SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: com.vitorpamplona.quartz.nip01Core.core.Event,
isLive: Boolean,
relay: com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl,
@@ -276,7 +276,7 @@ class Nip01ComplianceTest : RelayClientTest() {
"sub-A",
mapOf(defaultRelayUrl to listOf(Filter(kinds = listOf(1)))),
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -297,7 +297,7 @@ class Nip01ComplianceTest : RelayClientTest() {
"sub-B",
mapOf(defaultRelayUrl to listOf(Filter(kinds = listOf(4)))),
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -385,7 +385,7 @@ class Nip01ComplianceTest : RelayClientTest() {
"live-1",
mapOf(defaultRelayUrl to listOf(Filter(kinds = listOf(1)))),
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -424,7 +424,7 @@ class Nip01ComplianceTest : RelayClientTest() {
"live-2",
mapOf(defaultRelayUrl to listOf(Filter(kinds = listOf(1)))),
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -465,7 +465,7 @@ class Nip01ComplianceTest : RelayClientTest() {
"eph-1",
mapOf(defaultRelayUrl to listOf(Filter(kinds = listOf(20_001)))),
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -534,7 +534,7 @@ class Nip01ComplianceTest : RelayClientTest() {
relayB to listOf(Filter(kinds = listOf(1))),
),
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -80,7 +80,7 @@ class Nip09DeletionTest {
subId,
mapOf(relayUrl to listOf(filter)),
object : com.vitorpamplona.quartz.nip01Core.relay.client.reqs.SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -96,7 +96,7 @@ class Nip77NegentropyTest {
compression: Boolean,
) {}
override fun onMessage(text: String) {
override suspend fun onMessage(text: String) {
incoming.trySend(text)
}
@@ -99,7 +99,7 @@ class MirrorWorkerTrustOriginTest {
override fun connect() {
connected = true
out.onOpen(0, false)
out.onMessage(frame)
kotlinx.coroutines.runBlocking { out.onMessage(frame) }
}
override fun disconnect() {
@@ -258,7 +258,7 @@ class LoadBenchmark {
"fanout-$i",
mapOf(relayUrl to listOf(Filter(kinds = listOf(1)))),
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: com.vitorpamplona.quartz.nip01Core.core.Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -436,7 +436,7 @@ class LoadBenchmark {
"fanout-$i",
mapOf(relayUrl to listOf(Filter(kinds = listOf(1)))),
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: com.vitorpamplona.quartz.nip01Core.core.Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -568,7 +568,7 @@ class LoadBenchmark {
),
),
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: com.vitorpamplona.quartz.nip01Core.core.Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -72,7 +72,7 @@ suspend fun NostrClient.collectUntilEoseMulti(
subId,
mapOf(relay to filters),
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -1489,7 +1489,7 @@ class GrapeRankCrawler(
val lastEvt = AtomicLong(-1)
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -329,7 +329,7 @@ class NostrClient(
listeners.forEach { it.onSent(relay, cmdStr, cmd, success) }
}
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -137,7 +137,7 @@ class AdaptiveRelayLimiter(
if (wait > 0) delay(wait)
}
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -37,7 +37,7 @@ class EventCollector(
) {
private val clientListener =
object : RelayConnectionListener {
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -59,7 +59,7 @@ suspend fun INostrClient.count(
val listener =
object : RelayConnectionListener {
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -117,7 +117,7 @@ suspend fun INostrClient.count(
val listener =
object : RelayConnectionListener {
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -94,7 +94,7 @@ suspend fun INostrClient.fetchAllPages(
filters: List<Filter>,
idleTimeoutMs: Long = 30_000L,
onNewPage: ((Long) -> Unit)? = null,
onEvent: (Event) -> Unit,
onEvent: suspend (Event) -> Unit,
): Int {
var until: Long? = null
var totalEvents = 0
@@ -172,7 +172,7 @@ suspend fun INostrClient.fetchAllPages(
try {
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -320,7 +320,7 @@ suspend fun INostrClient.fetchAllPages(
filters: List<Filter>,
idleTimeoutMs: Long = 30_000L,
onNewPage: ((Long) -> Unit)? = null,
onEvent: (Event) -> Unit,
onEvent: suspend (Event) -> Unit,
): Int =
fetchAllPages(
relay = RelayUrlNormalizer.normalize(relay),
@@ -102,7 +102,7 @@ suspend fun INostrClient.fetchAllWithHooks(
val doneReasons = HashMap<NormalizedRelayUrl, String>()
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -96,7 +96,7 @@ suspend fun INostrClient.fetchFirst(
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -83,7 +83,7 @@ suspend fun negentropySyncFanOut(
idleTimeoutMs: Long = 120_000L,
reconcileConcurrency: Int = 2,
onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null,
onEvent: (Event) -> Unit,
onEvent: suspend (Event) -> Unit,
): NegentropyFanOutResult {
require(clients.isNotEmpty()) { "at least one client is required" }
@@ -163,7 +163,7 @@ suspend fun INostrClient.negentropySync(
idBufferBatches: Int = maxConcurrentReqs * 4,
localEntries: List<IdAndTime> = emptyList(),
onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null,
onEvent: (Event) -> Unit,
onEvent: suspend (Event) -> Unit,
): NegentropySyncResult {
val need = AtomicInt(0)
val windows = AtomicInt(0)
@@ -242,7 +242,7 @@ suspend fun INostrClient.negentropySync(
idBufferBatches: Int = maxConcurrentReqs * 4,
localEntries: List<IdAndTime> = emptyList(),
onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null,
onEvent: (Event) -> Unit,
onEvent: suspend (Event) -> Unit,
): NegentropySyncResult =
negentropySync(
relay = RelayUrlNormalizer.normalize(relay),
@@ -309,14 +309,14 @@ suspend fun INostrClient.negentropySyncOrFetch(
idBufferBatches: Int = maxConcurrentReqs * 4,
localEntries: List<IdAndTime> = emptyList(),
onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null,
onEvent: (Event) -> Unit,
onEvent: suspend (Event) -> Unit,
): NegentropyOrFetchResult {
val seen = HashSet<HexKey>()
var delivered = 0
// Shared dedup + cap across both phases. Returns true if the event was new and
// delivered. Both phases run sequentially, so no concurrent access.
fun accept(event: Event): Boolean {
suspend fun accept(event: Event): Boolean {
if ((maxEvents <= 0 || delivered < maxEvents) && seen.add(event.id)) {
delivered++
onEvent(event)
@@ -364,7 +364,7 @@ suspend fun INostrClient.negentropySyncOrFetch(
idBufferBatches: Int = maxConcurrentReqs * 4,
localEntries: List<IdAndTime> = emptyList(),
onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null,
onEvent: (Event) -> Unit,
onEvent: suspend (Event) -> Unit,
): NegentropyOrFetchResult =
negentropySyncOrFetch(
relay = RelayUrlNormalizer.normalize(relay),
@@ -876,7 +876,7 @@ private suspend fun INostrClient.reconcileStreaming(
if (relay.url == targetUrl) clock.bump()
}
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -1103,7 +1103,7 @@ internal suspend fun INostrClient.fetchByIds(
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -132,7 +132,7 @@ suspend fun INostrClient.publishAndCollectResults(
}
}
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -37,7 +37,7 @@ class RelayInsertConfirmationCollector(
) {
private val clientListener =
object : RelayConnectionListener {
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -50,7 +50,7 @@ class RelayLogger(
private val clientListener =
object : RelayConnectionListener {
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -40,7 +40,7 @@ class RelayNotifier(
private val clientListener =
object : RelayConnectionListener {
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -102,7 +102,7 @@ class RelayAuthenticator(
private val clientListener =
object : RelayConnectionListener {
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -45,7 +45,7 @@ class RelayActiveCountStates(
queryStates.put(relay.url, CountQueryState())
}
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -76,7 +76,7 @@ class RelayLimitsTracker(
private val clientListener =
object : RelayConnectionListener {
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -48,7 +48,7 @@ open class RedirectConnectionListener(
listener.onSent(relay, cmdStr, cmd, success)
}
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -53,7 +53,7 @@ interface RelayConnectionListener {
/**
* New error
*/
fun onIncomingMessage(
suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -123,7 +123,7 @@ class PoolCounts {
}
}
fun onIncomingMessage(
suspend fun onIncomingMessage(
relay: IRelayClient,
msg: Message,
) {
@@ -178,7 +178,7 @@ class PoolEventOutbox {
}
}
fun onIncomingMessage(
suspend fun onIncomingMessage(
relay: NormalizedRelayUrl,
msg: Message,
) {
@@ -238,7 +238,7 @@ class PoolRequests(
/**
* When a new message is received by the relay, updates the sub
*/
fun onIncomingMessage(
suspend fun onIncomingMessage(
relay: IRelayClient,
msg: Message,
) {
@@ -248,7 +248,7 @@ class RelayPool(
listener.onDisconnected(relay)
}
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -34,7 +34,7 @@ class DynamicSubscription(
SubscriptionHandle {
val subId = RandomInstance.randomChars(10)
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -63,7 +63,7 @@ fun INostrClient.fetchAsFlow(
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -69,7 +69,7 @@ fun INostrClient.subscribeAsFlow(
val listener =
object : SubscriptionListener {
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -46,7 +46,7 @@ class RelayActiveRequestStates(
subStates[relay.url] = RequestSubscriptionState()
}
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -35,7 +35,7 @@ class StaticSubscription(
SubscriptionHandle {
val subId = RandomInstance.randomChars(10)
override fun onEvent(
override suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -30,7 +30,7 @@ interface SubscriptionListener {
forFilters: List<Filter>?,
) {}
fun onEvent(
suspend fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
@@ -37,7 +37,7 @@ class RelayReqStats(
private val clientListener =
object : RelayConnectionListener {
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -156,7 +156,7 @@ open class BasicRelayClient(
listener.onConnected(this@BasicRelayClient, pingMillis, compression)
}
override fun onMessage(text: String) {
override suspend fun onMessage(text: String) {
try {
val msg = decoder.decode(text)
listener.onIncomingMessage(this@BasicRelayClient, text, msg)
@@ -66,7 +66,7 @@ class StandaloneRelayClient(
syncFilters()
}
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
@@ -83,7 +83,7 @@ class RelayStats(
get(relay.url).addBytesSent(cmdStr.bytesUsedInMemory())
}
override fun onIncomingMessage(
override suspend fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,

Some files were not shown because too many files have changed in this diff Show More