From da62bede43c9682835da837b8ffaa07198d82040 Mon Sep 17 00:00:00 2001 From: Kgothatso Ngako Date: Sun, 29 Mar 2026 20:43:27 +0200 Subject: [PATCH] Add NostrEventBroadcaster.kt --- .../compose/network/NostrEventBroadcaster.kt | 129 ++++++++++++++++++ .../ui/view/model/NavigationViewModel.kt | 115 ++++------------ 2 files changed, 153 insertions(+), 91 deletions(-) create mode 100644 composeApp/src/commonMain/kotlin/ac/aux/compose/network/NostrEventBroadcaster.kt diff --git a/composeApp/src/commonMain/kotlin/ac/aux/compose/network/NostrEventBroadcaster.kt b/composeApp/src/commonMain/kotlin/ac/aux/compose/network/NostrEventBroadcaster.kt new file mode 100644 index 00000000..8ecc006c --- /dev/null +++ b/composeApp/src/commonMain/kotlin/ac/aux/compose/network/NostrEventBroadcaster.kt @@ -0,0 +1,129 @@ +package ac.aux.compose.network + +import ac.aux.compose.database.model.BroadcastNostrEventReceipt +import ac.aux.compose.database.model.BroadcastNostrEventRequest +import ac.aux.compose.database.model.intermdiate.LocalBroadcastNostrEventRequest +import co.touchlab.kermit.Logger +import com.vitorpamplona.quartz.nip01Core.core.Event +import com.vitorpamplona.quartz.utils.text +import io.ktor.client.HttpClient +import io.ktor.client.plugins.websocket.WebSockets +import io.ktor.client.plugins.websocket.webSocketSession +import io.ktor.serialization.kotlinx.KotlinxWebsocketSerializationConverter +import io.ktor.websocket.Frame +import io.ktor.websocket.close +import io.ktor.websocket.readText +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Job +import kotlinx.coroutines.flow.consumeAsFlow +import kotlinx.coroutines.launch +import kotlinx.serialization.json.Json +import kotlinx.serialization.json.JsonArray + +class NostrEventBroadcaster( + val scope: CoroutineScope, + +) { + private val logger = Logger.withTag("NostrEventBroadcaster") + + private val broadcastingJobs = mutableMapOf() + + val httpClient = HttpClient() { + install(WebSockets) { + contentConverter = KotlinxWebsocketSerializationConverter(Json { + isLenient = true + ignoreUnknownKeys = true + }) + pingIntervalMillis = 20_000 + } + } + + fun broadcastEvent( + localBroadcastNostrEventRequest: LocalBroadcastNostrEventRequest, + onBroadcastRequestProcessed: (BroadcastNostrEventRequest) -> Unit, + onBroadcastReceipt: (BroadcastNostrEventReceipt) -> Unit + ) { + val event: Event = localBroadcastNostrEventRequest.nostrEvent.let { nostrEvent -> + Event( + id = nostrEvent.id, + kind = nostrEvent.kind, + pubKey = nostrEvent.pubKey, + content = nostrEvent.content, + tags = nostrEvent.tags, + sig = nostrEvent.sig, + createdAt = nostrEvent.createdAt.toEpochMilliseconds() + ) + } + val eventJson = event.toJson() + logger.d("Broadcasting (${localBroadcastNostrEventRequest.broadcastNostrEventRequest.relayURL}): $eventJson") + + scope.launch { + val webSocketSession = httpClient.webSocketSession( + urlString = localBroadcastNostrEventRequest.broadcastNostrEventRequest.relayURL + ) + + broadcastingJobs[localBroadcastNostrEventRequest.broadcastNostrEventRequest.relayURL] = scope.launch { + webSocketSession.send( + frame = Frame.Text("[\"EVENT\",${eventJson}]") + ) + + onBroadcastRequestProcessed.invoke(localBroadcastNostrEventRequest.broadcastNostrEventRequest) + + webSocketSession.incoming.consumeAsFlow().collect { frame -> + when(frame) { + is Frame.Text -> { + frame.readText().let { text -> + try { + logger.d("WebSocket Read Text: $text") + val result = Json.decodeFromString(text) + + if (result.firstOrNull()?.text == "OK") { + result.getOrNull(1)?.text?.let { eventId -> + if (eventId == localBroadcastNostrEventRequest.broadcastNostrEventRequest.nostrEventId) { + + onBroadcastReceipt.invoke( + BroadcastNostrEventReceipt( + nostrEventId = localBroadcastNostrEventRequest.broadcastNostrEventRequest.nostrEventId, + unsignedNostrEventId = localBroadcastNostrEventRequest.broadcastNostrEventRequest.unsignedNostrEventId, + messages = result.getOrNull(3)?.text, + isAccepted = result.getOrNull(2)?.text == "true", + relayURL = localBroadcastNostrEventRequest.broadcastNostrEventRequest.relayURL + ) + ) + } + } + } + } catch (e: Throwable) { + logger.e("Error processing receipt: $text", e) + } + } + webSocketSession.close() + broadcastingJobs.remove(localBroadcastNostrEventRequest.broadcastNostrEventRequest.relayURL)?.let { job -> + job.cancel() + } + } + is Frame.Close -> { + logger.d("Close $frame") + } + is Frame.Ping -> { + logger.d("Ping: $frame") + } + is Frame.Pong -> { + logger.d("Pong: $frame") + } + is Frame.Binary -> { + logger.d("Binary Message: $frame") + } + else -> { + logger.d("Unsupported frame: $frame") + } + } + } + + logger.i("Websocket exit") + } + + broadcastingJobs[localBroadcastNostrEventRequest.broadcastNostrEventRequest.relayURL]?.join() + } + } +} \ No newline at end of file diff --git a/composeApp/src/commonMain/kotlin/ac/aux/compose/ui/view/model/NavigationViewModel.kt b/composeApp/src/commonMain/kotlin/ac/aux/compose/ui/view/model/NavigationViewModel.kt index fcd2a1a1..5d86d7a6 100644 --- a/composeApp/src/commonMain/kotlin/ac/aux/compose/ui/view/model/NavigationViewModel.kt +++ b/composeApp/src/commonMain/kotlin/ac/aux/compose/ui/view/model/NavigationViewModel.kt @@ -3,6 +3,7 @@ package ac.aux.compose.ui.view.model import ac.aux.compose.database.model.BroadcastNostrEventReceipt import ac.aux.compose.database.model.NostrEvent import ac.aux.compose.managers.SeedManager +import ac.aux.compose.network.NostrEventBroadcaster import ac.aux.compose.nostr.Relays import ac.aux.compose.repository.NostrRepository import ac.aux.compose.ui.view.state.NavigationUIState @@ -61,7 +62,6 @@ class NavigationViewModel( } private val logger = Logger.withTag(TAG) - private val broadcastingJobs = mutableMapOf() val httpClient = HttpClient() { install(WebSockets) { @@ -73,6 +73,9 @@ class NavigationViewModel( } } + val nostrEventBroadcaster = NostrEventBroadcaster( + scope = scope + ) private val _navigationUIState = MutableStateFlow( initialNavigationUIState @@ -83,6 +86,7 @@ class NavigationViewModel( observeUnsignedNostrEvents() observePendingBroadcastNostrEventRequests() observeProfile() + observeSyncNostrEventRequests() } @@ -104,7 +108,6 @@ class NavigationViewModel( content = unsignedNostrEvent.content ) - nostrRepository.publishNostrEvent( unsignedNostrEvent, NostrEvent( @@ -125,6 +128,10 @@ class NavigationViewModel( } } + private fun observeSyncNostrEventRequests() { + + } + private fun observePendingBroadcastNostrEventRequests() { logger.i { "observePendingBroadcastNostrEventRequests" } @@ -132,98 +139,24 @@ class NavigationViewModel( nostrRepository.observePendingBroadcastNostrEventRequests().collect { localBroadcastNostrEventRequests -> localBroadcastNostrEventRequests.forEach { localBroadcastNostrEventRequest -> - val event: Event = localBroadcastNostrEventRequest.nostrEvent.let { nostrEvent -> - Event( - id = nostrEvent.id, - kind = nostrEvent.kind, - pubKey = nostrEvent.pubKey, - content = nostrEvent.content, - tags = nostrEvent.tags, - sig = nostrEvent.sig, - createdAt = nostrEvent.createdAt.toEpochMilliseconds() - ) - } - val eventJson = event.toJson() - logger.d("Broadcasting (${localBroadcastNostrEventRequest.broadcastNostrEventRequest.relayURL}): $eventJson") - - scope.launch { - - val webSocketSession = httpClient.webSocketSession( - urlString = localBroadcastNostrEventRequest.broadcastNostrEventRequest.relayURL - ) - - broadcastingJobs[localBroadcastNostrEventRequest.broadcastNostrEventRequest.relayURL] = scope.launch { - webSocketSession.send( - frame = Frame.Text("[\"EVENT\",${eventJson}]") - ) - - nostrRepository.broadcastProcessed( - localBroadcastNostrEventRequest.broadcastNostrEventRequest - ) - - webSocketSession.incoming.consumeAsFlow().collect { frame -> - when(frame) { - is Frame.Text -> { - frame.readText().let { text -> - try { - logger.d("WebSocket Read Text: $text") - val result = Json.decodeFromString(text) - - if (result.firstOrNull()?.text == "OK") { - result.getOrNull(1)?.text?.let { eventId -> - if (eventId == localBroadcastNostrEventRequest.broadcastNostrEventRequest.nostrEventId) { - nostrRepository.saveBroadcastReceipt( - BroadcastNostrEventReceipt( - nostrEventId = localBroadcastNostrEventRequest.broadcastNostrEventRequest.nostrEventId, - unsignedNostrEventId = localBroadcastNostrEventRequest.broadcastNostrEventRequest.unsignedNostrEventId, - messages = result.getOrNull(3)?.text, - isAccepted = result.getOrNull(2)?.text == "true", - relayURL = localBroadcastNostrEventRequest.broadcastNostrEventRequest.relayURL - ) - ) - } - } - } - } catch (e: Throwable) { - logger.e("Error processing receipt: $text", e) - } - } - webSocketSession.close() - broadcastingJobs.remove(localBroadcastNostrEventRequest.broadcastNostrEventRequest.relayURL)?.let { job -> - job.cancel() - } - } - is Frame.Close -> { - logger.d("Close $frame") - } - is Frame.Ping -> { - logger.d("Ping: $frame") - } - is Frame.Pong -> { - logger.d("Pong: $frame") - } - is Frame.Binary -> { - logger.d("Binary Message: $frame") - } - else -> { - logger.d("Unsupported frame: $frame") - } - } + nostrEventBroadcaster.broadcastEvent( + localBroadcastNostrEventRequest = localBroadcastNostrEventRequest, + onBroadcastRequestProcessed = { broadcastNostrEventRequest -> + scope.launch(Dispatchers.IO) { + nostrRepository.broadcastProcessed( + broadcastNostrEventRequest + ) + } + }, + onBroadcastReceipt = { broadcastNostrEventReceipt -> + scope.launch(Dispatchers.IO) { + nostrRepository.saveBroadcastReceipt( + broadcastNostrEventReceipt + ) } - - logger.i("Websocket exit") } - - broadcastingJobs[localBroadcastNostrEventRequest.broadcastNostrEventRequest.relayURL]?.join() -// nostrClient.send( -// event = event, -// relayList = setOf( -// RelayUrlNormalizer.normalize(localBroadcastNostrEventRequest.broadcastNostrEventRequest.relayURL) -// ) -// ) - } + ) } - } } }