Remove unused code
This commit is contained in:
@@ -1,204 +0,0 @@
|
||||
package ac.aux.compose.network
|
||||
|
||||
import co.touchlab.kermit.Logger
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebSocket
|
||||
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
|
||||
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.channels.SendChannel
|
||||
import kotlinx.coroutines.flow.consumeAsFlow
|
||||
import kotlinx.coroutines.isActive
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.runBlocking
|
||||
|
||||
class KTorHttpWebSocket(
|
||||
val url: NormalizedRelayUrl,
|
||||
val httpClientBuilder: (url: NormalizedRelayUrl) -> HttpClient,
|
||||
val ktorWebsocketListener: KtorWebsocketListener
|
||||
) : WebSocket {
|
||||
val scope = CoroutineScope(Dispatchers.IO)
|
||||
|
||||
val logger = Logger.withTag("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}")
|
||||
if (webSocketSession == null) return true
|
||||
|
||||
val activeHttpClient = httpClient ?: return true
|
||||
|
||||
val currentHttpClient = httpClientBuilder(url)
|
||||
|
||||
// TODO: Proxy...
|
||||
|
||||
// TODO: timeout checks
|
||||
|
||||
return false
|
||||
}
|
||||
|
||||
override fun connect() {
|
||||
logger.d("connect ${url.url}")
|
||||
runBlocking(
|
||||
scope.coroutineContext
|
||||
) {
|
||||
httpClient = httpClientBuilder(url)
|
||||
|
||||
webSocketSession = httpClient?.webSocketSession(
|
||||
urlString = url.url
|
||||
)
|
||||
|
||||
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")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
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")
|
||||
// }
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
override fun disconnect() {
|
||||
logger.d("disconnect ${url.url}")
|
||||
runBlocking(
|
||||
scope.coroutineContext
|
||||
) {
|
||||
webSocketSession?.close()
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
override fun send(msg: String): Boolean = try {
|
||||
logger.d("Send ${url.url}: $msg")
|
||||
runBlocking(
|
||||
scope.coroutineContext
|
||||
) {
|
||||
webSocketSession?.send(
|
||||
Frame.Text(msg)
|
||||
)
|
||||
|
||||
return@runBlocking true
|
||||
}
|
||||
} catch(e: Throwable) {
|
||||
logger.e("Failed to send $msg", e)
|
||||
return false
|
||||
}
|
||||
|
||||
class Builder(
|
||||
val ktorWebsocketListener: KtorWebsocketListener = KtorWebsocketListener(),
|
||||
val httpClientBuilder: (NormalizedRelayUrl) -> HttpClient,
|
||||
): WebsocketBuilder {
|
||||
val logger = Logger.withTag("KtorWebsocketBuilder")
|
||||
|
||||
override fun build(
|
||||
url: NormalizedRelayUrl,
|
||||
out: WebSocketListener
|
||||
): WebSocket {
|
||||
ktorWebsocketListener.out = out
|
||||
logger.d("Build KTorHttpWebSocket: ${url.url}")
|
||||
|
||||
return KTorHttpWebSocket(
|
||||
url = url,
|
||||
httpClientBuilder = httpClientBuilder,
|
||||
ktorWebsocketListener = ktorWebsocketListener
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,30 +0,0 @@
|
||||
package ac.aux.compose.network
|
||||
|
||||
import co.touchlab.kermit.Logger
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebSocketListener
|
||||
|
||||
class KtorWebsocketListener: WebSocketListener {
|
||||
val logger = Logger.withTag("KtorWebsocketListener")
|
||||
|
||||
var out: WebSocketListener? = null
|
||||
|
||||
override fun onOpen(pingMillis: Int, compression: Boolean) {
|
||||
logger.d("onOpen: $pingMillis, $compression")
|
||||
out?.onOpen(pingMillis, compression)
|
||||
}
|
||||
|
||||
override fun onMessage(text: String) {
|
||||
logger.d("onMessage: $text")
|
||||
out?.onMessage(text)
|
||||
}
|
||||
|
||||
override fun onClosed(code: Int, reason: String) {
|
||||
logger.d("onClosed: $code, $reason")
|
||||
out?.onClosed(code, reason)
|
||||
}
|
||||
|
||||
override fun onFailure(t: Throwable, code: Int?, response: String?) {
|
||||
logger.e("onFailure: $code, $response", t)
|
||||
out?.onFailure(t, code, response)
|
||||
}
|
||||
}
|
||||
@@ -1,142 +0,0 @@
|
||||
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.nip01Core.core.OptimizedJsonMapper
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.EventCmd
|
||||
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.epochSeconds
|
||||
)
|
||||
}
|
||||
|
||||
scope.launch {
|
||||
try {
|
||||
val webSocketSession = httpClient.webSocketSession(
|
||||
urlString = localBroadcastNostrEventRequest.broadcastNostrEventRequest.relayURL
|
||||
)
|
||||
|
||||
broadcastingJobs[localBroadcastNostrEventRequest.broadcastNostrEventRequest.relayURL] = scope.launch {
|
||||
val message = OptimizedJsonMapper.toJson(
|
||||
EventCmd(
|
||||
event
|
||||
)
|
||||
)
|
||||
logger.d("Broadcasting (${localBroadcastNostrEventRequest.broadcastNostrEventRequest.relayURL}): $message")
|
||||
webSocketSession.send(
|
||||
frame = Frame.Text(message)
|
||||
)
|
||||
|
||||
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()
|
||||
} catch (e: Throwable) {
|
||||
logger.e("Error broadcasting event: ", e)
|
||||
// TODO: Save this a failure
|
||||
onBroadcastRequestProcessed.invoke(localBroadcastNostrEventRequest.broadcastNostrEventRequest)
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -3,7 +3,6 @@ package ac.aux.compose.ui.view.model
|
||||
import ac.aux.compose.database.model.NostrEvent
|
||||
import ac.aux.compose.managers.SeedManager
|
||||
import ac.aux.compose.network.dto.RelayDTO
|
||||
import ac.aux.compose.network.relays.RelayPool
|
||||
import ac.aux.compose.network.relays.RelayPool.Companion.PUBLISH_TIMEOUT
|
||||
import ac.aux.compose.network.relays.RelaysSocketManager
|
||||
import ac.aux.compose.network.sockets.NostrIncomingMessage
|
||||
@@ -37,7 +36,6 @@ import kotlinx.coroutines.flow.getAndUpdate
|
||||
import kotlinx.coroutines.flow.timeout
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlin.time.Duration.Companion.milliseconds
|
||||
import kotlin.time.Duration.Companion.seconds
|
||||
import kotlin.time.Instant
|
||||
|
||||
class NavigationViewModel(
|
||||
|
||||
Reference in New Issue
Block a user