From 9e8dec41ca8f88b961abfd80b4455f78faf77424 Mon Sep 17 00:00:00 2001 From: Kgothatso Ngako Date: Sat, 28 Mar 2026 23:13:42 +0200 Subject: [PATCH] Send nostr events via websockets --- composeApp/build.gradle.kts | 1 + .../repository/DatabaseNostrRepository.kt | 7 + .../aux/compose/network/KTorHttpWebSocket.kt | 143 ++++++++++++------ .../compose/network/KtorWebsocketListener.kt | 2 +- .../aux/compose/repository/NostrRepository.kt | 7 + .../ui/view/model/NavigationViewModel.kt | 102 +++++++++++-- 6 files changed, 201 insertions(+), 61 deletions(-) diff --git a/composeApp/build.gradle.kts b/composeApp/build.gradle.kts index 73b0ccd4..a21d65c5 100644 --- a/composeApp/build.gradle.kts +++ b/composeApp/build.gradle.kts @@ -70,6 +70,7 @@ kotlin { implementation(libs.ktor.client.core) implementation("io.ktor:ktor-client-websockets:3.4.1") + implementation("io.ktor:ktor-serialization-kotlinx-json:3.4.1") } commonTest.dependencies { implementation(libs.kotlin.test) diff --git a/composeApp/src/commonMain/kotlin/ac/aux/compose/database/repository/DatabaseNostrRepository.kt b/composeApp/src/commonMain/kotlin/ac/aux/compose/database/repository/DatabaseNostrRepository.kt index 4c8d7e3a..f76bf870 100644 --- a/composeApp/src/commonMain/kotlin/ac/aux/compose/database/repository/DatabaseNostrRepository.kt +++ b/composeApp/src/commonMain/kotlin/ac/aux/compose/database/repository/DatabaseNostrRepository.kt @@ -1,6 +1,7 @@ package ac.aux.compose.database.repository import ac.aux.compose.database.AuxDatabase +import ac.aux.compose.database.model.BroadcastNostrEventReceipt import ac.aux.compose.database.model.BroadcastNostrEventRequest import ac.aux.compose.database.model.NostrEvent import ac.aux.compose.database.model.Profile @@ -104,4 +105,10 @@ class DatabaseNostrRepository( logger.d("Processed: $broadcastNostrEventRequest") } + override suspend fun saveBroadcastReceipt(broadcastNostrEventReceipt: BroadcastNostrEventReceipt) { + database.broadcastNostrEventReceiptDao().upsert( + broadcastNostrEventReceipt + ) + } + } \ No newline at end of file diff --git a/composeApp/src/commonMain/kotlin/ac/aux/compose/network/KTorHttpWebSocket.kt b/composeApp/src/commonMain/kotlin/ac/aux/compose/network/KTorHttpWebSocket.kt index e4376821..34b455dd 100644 --- a/composeApp/src/commonMain/kotlin/ac/aux/compose/network/KTorHttpWebSocket.kt +++ b/composeApp/src/commonMain/kotlin/ac/aux/compose/network/KTorHttpWebSocket.kt @@ -7,6 +7,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebSocketListener import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebsocketBuilder import io.ktor.client.HttpClient import io.ktor.client.plugins.websocket.DefaultClientWebSocketSession +import io.ktor.client.plugins.websocket.webSocket import io.ktor.client.plugins.websocket.webSocketSession import io.ktor.websocket.Frame import io.ktor.websocket.close @@ -15,6 +16,8 @@ import io.ktor.websocket.readText import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.IO +import kotlinx.coroutines.channels.SendChannel +import kotlinx.coroutines.flow.consumeAsFlow import kotlinx.coroutines.isActive import kotlinx.coroutines.launch import kotlinx.coroutines.runBlocking @@ -22,7 +25,7 @@ import kotlinx.coroutines.runBlocking class KTorHttpWebSocket( val url: NormalizedRelayUrl, val httpClientBuilder: (url: NormalizedRelayUrl) -> HttpClient, - val out: WebSocketListener + val ktorWebsocketListener: KtorWebsocketListener ) : WebSocket { val scope = CoroutineScope(Dispatchers.IO) @@ -31,6 +34,7 @@ class KTorHttpWebSocket( private var httpClient: HttpClient? = null private var webSocketSession: DefaultClientWebSocketSession? = null + private var outgoingChannel: SendChannel? = null override fun needsReconnect(): Boolean { logger.d("needsReconnect ${url.url}") @@ -53,62 +57,100 @@ class KTorHttpWebSocket( scope.coroutineContext ) { httpClient = httpClientBuilder(url) + webSocketSession = httpClient?.webSocketSession( urlString = url.url - ) { - logger.d("webSocketSession") - } + ) - webSocketSession?.let { socketSession -> - - logger.d("isOpen : ${socketSession.isActive}") - out.onOpen( - pingMillis = 1, - compression = false - ) - - val incomingFrameJob = scope.launch { - try { - for (frame in socketSession.incoming) { - when (frame) { - is Frame.Text -> { - frame.readText().let { - logger.d("WebSocket Read Text: $it") - out.onMessage(it) - } - } - is Frame.Close -> { - logger.d("Close $frame") - val reason = frame.readReason() - out.onClosed( - code = reason?.code?.toInt() ?: 0, - reason = reason?.knownReason.toString() - ) - } - 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") - } + val job = scope.launch { + webSocketSession?.incoming?.consumeAsFlow()?.collect { frame -> + when(frame) { + is Frame.Text -> { + frame.readText().let { + logger.d("WebSocket Read Text: $it") } + webSocketSession?.close() + } + 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") } - } catch (e: Throwable) { - logger.e("Websocket error: ", e) - } finally { - logger.d("Closed for real... hopefully we know the reason") } } - incomingFrameJob.join() - logger.d("webSocketSession closed") + + logger.i("Websocket exit") } + + job.join() + ktorWebsocketListener.onClosed(0, "Done") +// webSocketSession = httpClient?.webSocketSession( +// urlString = url.url +// ) { +// logger.d("webSocketSession") +// } +// +// webSocketSession?.let { socketSession -> +// +// logger.d("isOpen : ${socketSession.isActive}") +// out.onOpen( +// pingMillis = 1, +// compression = false +// ) +// +// val incomingFrameJob = scope.launch { +// try { +// for (frame in socketSession.incoming) { +// when (frame) { +// is Frame.Text -> { +// frame.readText().let { +// logger.d("WebSocket Read Text: $it") +// out.onMessage(it) +// } +// } +// is Frame.Close -> { +// logger.d("Close $frame") +// val reason = frame.readReason() +// out.onClosed( +// code = reason?.code?.toInt() ?: 0, +// reason = reason?.knownReason.toString() +// ) +// } +// 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") +// } +// } +// } +// } catch (e: Throwable) { +// logger.e("Websocket error: ", e) +// } finally { +// logger.d("Closed for real... hopefully we know the reason") +// } +// } +// +// incomingFrameJob.join() +// logger.d("webSocketSession closed") +// } } } @@ -140,7 +182,7 @@ class KTorHttpWebSocket( } class Builder( - ktorWebsocketListener: KtorWebsocketListener = KtorWebsocketListener(), + val ktorWebsocketListener: KtorWebsocketListener = KtorWebsocketListener(), val httpClientBuilder: (NormalizedRelayUrl) -> HttpClient, ): WebsocketBuilder { val logger = Logger.withTag("KtorWebsocketBuilder") @@ -149,12 +191,13 @@ class KTorHttpWebSocket( url: NormalizedRelayUrl, out: WebSocketListener ): WebSocket { + ktorWebsocketListener.out = out logger.d("Build KTorHttpWebSocket: ${url.url}") return KTorHttpWebSocket( url = url, httpClientBuilder = httpClientBuilder, - out = out + ktorWebsocketListener = ktorWebsocketListener ) } } diff --git a/composeApp/src/commonMain/kotlin/ac/aux/compose/network/KtorWebsocketListener.kt b/composeApp/src/commonMain/kotlin/ac/aux/compose/network/KtorWebsocketListener.kt index 2e87b49c..e288c1ca 100644 --- a/composeApp/src/commonMain/kotlin/ac/aux/compose/network/KtorWebsocketListener.kt +++ b/composeApp/src/commonMain/kotlin/ac/aux/compose/network/KtorWebsocketListener.kt @@ -6,7 +6,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebSocketListener class KtorWebsocketListener: WebSocketListener { val logger = Logger.withTag("KtorWebsocketListener") - val out: WebSocketListener? = null + var out: WebSocketListener? = null override fun onOpen(pingMillis: Int, compression: Boolean) { logger.d("onOpen: $pingMillis, $compression") diff --git a/composeApp/src/commonMain/kotlin/ac/aux/compose/repository/NostrRepository.kt b/composeApp/src/commonMain/kotlin/ac/aux/compose/repository/NostrRepository.kt index ce8c8b11..cfa518b5 100644 --- a/composeApp/src/commonMain/kotlin/ac/aux/compose/repository/NostrRepository.kt +++ b/composeApp/src/commonMain/kotlin/ac/aux/compose/repository/NostrRepository.kt @@ -1,5 +1,6 @@ package ac.aux.compose.repository +import ac.aux.compose.database.model.BroadcastNostrEventReceipt import ac.aux.compose.database.model.BroadcastNostrEventRequest import ac.aux.compose.database.model.NostrEvent import ac.aux.compose.database.model.Profile @@ -38,6 +39,8 @@ interface NostrRepository { suspend fun broadcastProcessed(broadcastNostrEventRequest: BroadcastNostrEventRequest) + suspend fun saveBroadcastReceipt(broadcastNostrEventReceipt: BroadcastNostrEventReceipt) + companion object { val NO_OP_NOSTR_REPOSITORY = object : NostrRepository { override suspend fun observeProfile(publicKey: HexKey): Flow { @@ -79,6 +82,10 @@ interface NostrRepository { override suspend fun broadcastProcessed(broadcastNostrEventRequest: BroadcastNostrEventRequest) { TODO("Not yet implemented") } + + override suspend fun saveBroadcastReceipt(broadcastNostrEventReceipt: BroadcastNostrEventReceipt) { + TODO("Not yet implemented") + } } } } \ 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 88cd31eb..426ddd15 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 @@ -1,5 +1,6 @@ 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.KTorHttpWebSocket @@ -18,16 +19,30 @@ import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebsocketBuilder import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerSync +import com.vitorpamplona.quartz.utils.text import io.ktor.client.HttpClient import io.ktor.client.plugins.websocket.WebSockets +import io.ktor.client.plugins.websocket.sendSerialized +import io.ktor.client.plugins.websocket.webSocket +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.readReason +import io.ktor.websocket.readText import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.IO import kotlinx.coroutines.delay import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.asStateFlow +import kotlinx.coroutines.flow.consumeAsFlow import kotlinx.coroutines.flow.getAndUpdate import kotlinx.coroutines.launch +import kotlinx.serialization.json.Json +import kotlinx.serialization.json.JsonArray +import kotlinx.serialization.json.jsonObject +import kotlin.math.log import kotlin.time.Instant class NavigationViewModel( @@ -57,9 +72,12 @@ class NavigationViewModel( private val logger = Logger.withTag(TAG) - val httpClient = HttpClient() { install(WebSockets) { + contentConverter = KotlinxWebsocketSerializationConverter(Json { + isLenient = true + ignoreUnknownKeys = true + }) pingIntervalMillis = 20_000 } } @@ -146,20 +164,84 @@ class NavigationViewModel( createdAt = nostrEvent.createdAt.toEpochMilliseconds() ) } - logger.d("Broadcasting (${localBroadcastNostrEventRequest.broadcastNostrEventRequest.relayURL}): ${event.toJson()}") + val eventJson = event.toJson() + logger.d("Broadcasting (${localBroadcastNostrEventRequest.broadcastNostrEventRequest.relayURL}): $eventJson") scope.launch { - nostrClient.send( - event = event, - relayList = setOf( - RelayUrlNormalizer.normalize(localBroadcastNostrEventRequest.broadcastNostrEventRequest.relayURL) - ) + val webSocketSession = httpClient.webSocketSession( + urlString = localBroadcastNostrEventRequest.broadcastNostrEventRequest.relayURL ) -// nostrRepository.broadcastProcessed( -// localBroadcastNostrEventRequest.broadcastNostrEventRequest -// ) + val job = 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 -> + // TODO: receipt + 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() + } + 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") + } + + job.join() +// nostrClient.send( +// event = event, +// relayList = setOf( +// RelayUrlNormalizer.normalize(localBroadcastNostrEventRequest.broadcastNostrEventRequest.relayURL) +// ) +// ) + + } }