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)
+// )
+// )
+
+
}
}