Add NostrEventBroadcaster.kt

This commit is contained in:
Kgothatso Ngako
2026-03-29 20:43:27 +02:00
parent 228ddbbea1
commit da62bede43
2 changed files with 153 additions and 91 deletions

View File

@@ -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<String, Job>()
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<JsonArray>(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()
}
}
}

View File

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