Merge pull request #3863 from vitorpamplona/feat/suspend-subscription-onevent

Suspend the incoming-message chain down to SubscriptionListener.onEvent
This commit is contained in:
Vitor Pamplona
2026-08-05 14:35:03 -04:00
committed by GitHub
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,
@@ -701,7 +701,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