Send nostr events via websockets

This commit is contained in:
Kgothatso Ngako
2026-03-28 23:13:42 +02:00
parent e4d253cbd7
commit 9e8dec41ca
6 changed files with 201 additions and 61 deletions

View File

@@ -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
)
}
}

View File

@@ -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<Frame>? = 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
)
}
}

View File

@@ -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")

View File

@@ -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<LocalProfile?> {
@@ -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")
}
}
}
}

View File

@@ -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<JsonArray>(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)
// )
// )
}
}