Broadcast the things
This commit is contained in:
@@ -1,6 +1,9 @@
|
||||
<?xml version="1.0" encoding="utf-8"?>
|
||||
<manifest xmlns:android="http://schemas.android.com/apk/res/android">
|
||||
|
||||
<!-- To connect with relays -->
|
||||
<uses-permission android:name="android.permission.INTERNET"/>
|
||||
|
||||
<application
|
||||
android:allowBackup="true"
|
||||
android:icon="@mipmap/ic_launcher"
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package ac.aux.compose.database.dao
|
||||
|
||||
import ac.aux.compose.database.model.BroadcastNostrEventRequest
|
||||
import ac.aux.compose.database.model.intermdiate.LocalBroadcastNostrEventRequest
|
||||
import androidx.room.Dao
|
||||
import androidx.room.Query
|
||||
import androidx.room.Upsert
|
||||
@@ -11,11 +12,11 @@ interface BroadcastNostrEventRequestDao {
|
||||
@Query("SELECT * FROM BroadcastNostrEventRequest")
|
||||
fun getAllBroadcastNostrEventRequests(): List<BroadcastNostrEventRequest>
|
||||
|
||||
@Query("SELECT * FROM BroadcastNostrEventRequest WHERE id = :id")
|
||||
fun getBroadcastNostrEventRequestById(id: Long): Flow<BroadcastNostrEventRequest?>
|
||||
@Query("SELECT * FROM BroadcastNostrEventRequest WHERE status = :status")
|
||||
fun observeBroadcastNostrEventRequestByStatus(status: String): Flow<LocalBroadcastNostrEventRequest?>
|
||||
|
||||
@Query("SELECT * FROM BroadcastNostrEventRequest WHERE nostrEventId = :nostrEventId")
|
||||
fun getBroadcastNostrEventRequestByNostrEventId(nostrEventId: String): Flow<BroadcastNostrEventRequest?>
|
||||
fun observeBroadcastNostrEventRequestByNostrEventId(nostrEventId: String): Flow<BroadcastNostrEventRequest?>
|
||||
|
||||
@Upsert
|
||||
suspend fun upsert(broadcastNostrEventRequest: BroadcastNostrEventRequest)
|
||||
|
||||
@@ -0,0 +1,19 @@
|
||||
package ac.aux.compose.database.model.intermdiate
|
||||
|
||||
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
|
||||
import ac.aux.compose.database.model.UnsignedNostrEvent
|
||||
import androidx.room.Embedded
|
||||
import androidx.room.Relation
|
||||
|
||||
data class LocalBroadcastNostrEventRequest(
|
||||
@Embedded val broadcastNostrEventRequest: BroadcastNostrEventRequest,
|
||||
|
||||
@Relation(
|
||||
parentColumn = "nostrEventId",
|
||||
entityColumn = "id"
|
||||
)
|
||||
val nostrEvent: NostrEvent,
|
||||
)
|
||||
@@ -1,9 +1,11 @@
|
||||
package ac.aux.compose.database.repository
|
||||
|
||||
import ac.aux.compose.database.AuxDatabase
|
||||
import ac.aux.compose.database.model.BroadcastNostrEventRequest
|
||||
import ac.aux.compose.database.model.NostrEvent
|
||||
import ac.aux.compose.database.model.Profile
|
||||
import ac.aux.compose.database.model.UnsignedNostrEvent
|
||||
import ac.aux.compose.database.model.intermdiate.LocalBroadcastNostrEventRequest
|
||||
import ac.aux.compose.database.model.intermdiate.LocalProfile
|
||||
import ac.aux.compose.repository.NostrRepository
|
||||
import co.touchlab.kermit.Logger
|
||||
@@ -30,6 +32,10 @@ class DatabaseNostrRepository(
|
||||
return database.unsignedNostrEventDao().observeUnsignedNostrEvents(publicKey)
|
||||
}
|
||||
|
||||
override suspend fun observePendingBroadcastNostrEventRequests(): Flow<LocalBroadcastNostrEventRequest?> {
|
||||
return database.broadcastNostrEventRequestDao().observeBroadcastNostrEventRequestByStatus("pending")
|
||||
}
|
||||
|
||||
override suspend fun createNewProfile(
|
||||
publicKey: HexKey,
|
||||
name: String,
|
||||
@@ -83,4 +89,14 @@ class DatabaseNostrRepository(
|
||||
)
|
||||
}
|
||||
|
||||
override suspend fun broadcastProcessed(broadcastNostrEventRequest: BroadcastNostrEventRequest) {
|
||||
logger.d("Update local reference: $broadcastNostrEventRequest")
|
||||
database.broadcastNostrEventRequestDao().upsert(
|
||||
broadcastNostrEventRequest.copy(
|
||||
status = "sent"
|
||||
)
|
||||
)
|
||||
logger.d("Processed: $broadcastNostrEventRequest")
|
||||
}
|
||||
|
||||
}
|
||||
@@ -0,0 +1,98 @@
|
||||
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.webSocketSession
|
||||
import io.ktor.websocket.Frame
|
||||
import io.ktor.websocket.close
|
||||
import io.ktor.websocket.readText
|
||||
import io.ktor.websocket.send
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.IO
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.runBlocking
|
||||
|
||||
class KTorHttpWebSocket(
|
||||
val url: NormalizedRelayUrl,
|
||||
val httpClientBuilder: (url: NormalizedRelayUrl) -> HttpClient,
|
||||
val out: WebSocketListener
|
||||
) : WebSocket {
|
||||
val scope = CoroutineScope(Dispatchers.IO)
|
||||
|
||||
val logger = Logger.withTag("KTorHttpWebSocket")
|
||||
|
||||
private var httpClient: HttpClient? = null
|
||||
|
||||
private var webSocketSession: DefaultClientWebSocketSession? = null
|
||||
|
||||
override fun needsReconnect(): Boolean {
|
||||
TODO("Not yet implemented")
|
||||
}
|
||||
|
||||
override fun connect() {
|
||||
runBlocking {
|
||||
httpClient = httpClientBuilder(url)
|
||||
webSocketSession = httpClient?.webSocketSession(
|
||||
urlString = url.url
|
||||
) {
|
||||
|
||||
}
|
||||
|
||||
webSocketSession?.let { socketSession ->
|
||||
|
||||
val job = scope.launch {
|
||||
for (frame in socketSession.incoming) {
|
||||
val frameText = frame as? Frame.Text
|
||||
|
||||
frameText?.readText()?.let { out.onMessage(it) }
|
||||
}
|
||||
}
|
||||
|
||||
job.join()
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
override fun disconnect() {
|
||||
runBlocking {
|
||||
webSocketSession?.close()
|
||||
}
|
||||
}
|
||||
|
||||
override fun send(msg: String): Boolean = try {
|
||||
runBlocking {
|
||||
|
||||
webSocketSession?.send(
|
||||
Frame.Text(msg)
|
||||
)
|
||||
|
||||
return@runBlocking true
|
||||
}
|
||||
} catch(e: Throwable) {
|
||||
logger.e("Failed to send $msg", e)
|
||||
return false
|
||||
}
|
||||
|
||||
class Builder(
|
||||
val httpClientBuilder: (NormalizedRelayUrl) -> HttpClient
|
||||
): WebsocketBuilder {
|
||||
override fun build(
|
||||
url: NormalizedRelayUrl,
|
||||
out: WebSocketListener
|
||||
): WebSocket = KTorHttpWebSocket(
|
||||
url = url,
|
||||
httpClientBuilder = httpClientBuilder,
|
||||
out = out
|
||||
)
|
||||
}
|
||||
|
||||
}
|
||||
@@ -1,8 +1,10 @@
|
||||
package ac.aux.compose.repository
|
||||
|
||||
import ac.aux.compose.database.model.BroadcastNostrEventRequest
|
||||
import ac.aux.compose.database.model.NostrEvent
|
||||
import ac.aux.compose.database.model.Profile
|
||||
import ac.aux.compose.database.model.UnsignedNostrEvent
|
||||
import ac.aux.compose.database.model.intermdiate.LocalBroadcastNostrEventRequest
|
||||
import ac.aux.compose.database.model.intermdiate.LocalProfile
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
@@ -12,6 +14,8 @@ interface NostrRepository {
|
||||
|
||||
suspend fun observeUnsignedNostrEvents(publicKey: HexKey): Flow<UnsignedNostrEvent?>
|
||||
|
||||
suspend fun observePendingBroadcastNostrEventRequests(): Flow<LocalBroadcastNostrEventRequest?>
|
||||
|
||||
suspend fun createNewProfile(
|
||||
publicKey: HexKey,
|
||||
name: String,
|
||||
@@ -31,6 +35,8 @@ interface NostrRepository {
|
||||
relayURLs: List<String> = emptyList()
|
||||
)
|
||||
|
||||
suspend fun broadcastProcessed(broadcastNostrEventRequest: BroadcastNostrEventRequest)
|
||||
|
||||
companion object {
|
||||
val NO_OP_NOSTR_REPOSITORY = object : NostrRepository {
|
||||
override suspend fun observeProfile(publicKey: HexKey): Flow<LocalProfile?> {
|
||||
@@ -41,6 +47,10 @@ interface NostrRepository {
|
||||
TODO("Not yet implemented")
|
||||
}
|
||||
|
||||
override suspend fun observePendingBroadcastNostrEventRequests(): Flow<LocalBroadcastNostrEventRequest?> {
|
||||
TODO("Not yet implemented")
|
||||
}
|
||||
|
||||
override suspend fun createNewProfile(
|
||||
publicKey: HexKey,
|
||||
name: String,
|
||||
@@ -60,6 +70,10 @@ interface NostrRepository {
|
||||
override suspend fun publishNostrEvent(nostrEvent: NostrEvent, relayURLs: List<String>) {
|
||||
|
||||
}
|
||||
|
||||
override suspend fun broadcastProcessed(broadcastNostrEventRequest: BroadcastNostrEventRequest) {
|
||||
TODO("Not yet implemented")
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -41,12 +41,19 @@ import androidx.compose.ui.Modifier
|
||||
import androidx.compose.ui.graphics.Color
|
||||
import androidx.lifecycle.compose.LocalLifecycleOwner
|
||||
import androidx.lifecycle.viewmodel.compose.viewModel
|
||||
import co.touchlab.kermit.Logger
|
||||
import kotlinx.coroutines.CoroutineExceptionHandler
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.IO
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
|
||||
@Composable
|
||||
fun AuxNavHost(
|
||||
auxGlobal: AuxGlobal,
|
||||
navController: NavHostController
|
||||
) {
|
||||
val logger = Logger.withTag("AuxNavHost")
|
||||
val lifecycleOwner = LocalLifecycleOwner.current
|
||||
|
||||
val databaseManager = DatabaseManager(auxGlobal)
|
||||
@@ -55,10 +62,19 @@ fun AuxNavHost(
|
||||
database = databaseManager.auxDatabase
|
||||
)
|
||||
|
||||
val exceptionHandler =
|
||||
CoroutineExceptionHandler { _, throwable ->
|
||||
logger.e("Caught exception: ${throwable.message}", throwable)
|
||||
}
|
||||
|
||||
val applicationIOScope = CoroutineScope(Dispatchers.IO + SupervisorJob() + exceptionHandler)
|
||||
|
||||
|
||||
val navigationViewModel: NavigationViewModel = viewModel (
|
||||
factory = NavigationViewModel.factory(
|
||||
NavigationUIState.Loading,
|
||||
nostrRepository
|
||||
initialNavigationUIState = NavigationUIState.Loading,
|
||||
nostrRepository = nostrRepository,
|
||||
scope = applicationIOScope
|
||||
)
|
||||
)
|
||||
|
||||
@@ -164,9 +180,7 @@ fun AuxNavHost(
|
||||
}
|
||||
)
|
||||
}
|
||||
composable<CreateProfileRoute> { backStackEntry ->
|
||||
val route = backStackEntry.toRoute<CreateProfileRoute>()
|
||||
|
||||
composable<CreateProfileRoute> {
|
||||
CreateProfileScreen(
|
||||
onNavigateToSocialPreconditionRoute = {
|
||||
navController.navigate(
|
||||
|
||||
@@ -2,6 +2,7 @@ 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.KTorHttpWebSocket
|
||||
import ac.aux.compose.nostr.Relays
|
||||
import ac.aux.compose.repository.NostrRepository
|
||||
import ac.aux.compose.ui.view.state.NavigationUIState
|
||||
@@ -9,11 +10,18 @@ import androidx.lifecycle.ViewModel
|
||||
import androidx.lifecycle.ViewModelProvider
|
||||
import androidx.lifecycle.viewmodel.initializer
|
||||
import androidx.lifecycle.viewmodel.viewModelFactory
|
||||
import androidx.lifecycle.viewModelScope
|
||||
import co.touchlab.kermit.Logger
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.core.toHexKey
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
|
||||
import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerSync
|
||||
import io.ktor.client.HttpClient
|
||||
import io.ktor.client.plugins.websocket.WebSockets
|
||||
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
|
||||
@@ -23,19 +31,23 @@ import kotlin.time.Instant
|
||||
|
||||
class NavigationViewModel(
|
||||
initialNavigationUIState: NavigationUIState,
|
||||
val nostrRepository: NostrRepository
|
||||
val nostrRepository: NostrRepository,
|
||||
val scope: CoroutineScope,
|
||||
): ViewModel() {
|
||||
|
||||
companion object {
|
||||
private const val TAG = "NavigationViewModel"
|
||||
|
||||
fun factory(
|
||||
initialNavigationUIState: NavigationUIState,
|
||||
nostrRepository: NostrRepository
|
||||
nostrRepository: NostrRepository,
|
||||
scope: CoroutineScope,
|
||||
): ViewModelProvider.Factory = viewModelFactory {
|
||||
initializer {
|
||||
NavigationViewModel(
|
||||
initialNavigationUIState = initialNavigationUIState,
|
||||
nostrRepository = nostrRepository
|
||||
nostrRepository = nostrRepository,
|
||||
scope = scope
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -43,6 +55,24 @@ class NavigationViewModel(
|
||||
|
||||
private val logger = Logger.withTag(TAG)
|
||||
|
||||
|
||||
|
||||
val httpClient = HttpClient() {
|
||||
install(WebSockets) {
|
||||
pingIntervalMillis = 20_000
|
||||
}
|
||||
}
|
||||
|
||||
val websocketBuilder = KTorHttpWebSocket.Builder { url ->
|
||||
// TODO: Figure out if we need a tor client
|
||||
httpClient
|
||||
}
|
||||
val nostrClient: INostrClient = NostrClient(
|
||||
websocketBuilder = websocketBuilder,
|
||||
scope = scope
|
||||
)
|
||||
|
||||
|
||||
private val _navigationUIState = MutableStateFlow(
|
||||
initialNavigationUIState
|
||||
)
|
||||
@@ -50,18 +80,17 @@ class NavigationViewModel(
|
||||
|
||||
init {
|
||||
observeUnsignedNostrEvents()
|
||||
observePendingBroadcastNostrEventRequests()
|
||||
observeProfile()
|
||||
}
|
||||
|
||||
|
||||
|
||||
|
||||
private fun observeUnsignedNostrEvents() {
|
||||
val tempSigner = NostrSignerSync(
|
||||
SeedManager.activeKeyPair()
|
||||
)
|
||||
logger.i { "observeUnsignedNostrEvents" }
|
||||
viewModelScope.launch {
|
||||
scope.launch(Dispatchers.IO) {
|
||||
nostrRepository.observeUnsignedNostrEvents(
|
||||
publicKey = SeedManager.activePublicKey().toHexKey()
|
||||
).collect { unsignedNostrEventOrNull ->
|
||||
@@ -92,9 +121,45 @@ class NavigationViewModel(
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private fun observePendingBroadcastNostrEventRequests() {
|
||||
logger.i { "observePendingBroadcastNostrEventRequests" }
|
||||
|
||||
scope.launch(Dispatchers.IO) {
|
||||
|
||||
nostrRepository.observePendingBroadcastNostrEventRequests().collect {
|
||||
it?.let { 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()
|
||||
)
|
||||
}
|
||||
logger.d("Broadcasting: ${event.toJson()}")
|
||||
nostrClient.send(
|
||||
event = event,
|
||||
relayList = setOf(
|
||||
RelayUrlNormalizer.normalize(localBroadcastNostrEventRequest.broadcastNostrEventRequest.relayURL)
|
||||
)
|
||||
)
|
||||
logger.d("Update with result")
|
||||
|
||||
nostrRepository.broadcastProcessed(
|
||||
localBroadcastNostrEventRequest.broadcastNostrEventRequest
|
||||
)
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
}
|
||||
private fun observeProfile() {
|
||||
logger.i("observeProfile")
|
||||
viewModelScope.launch {
|
||||
scope.launch(Dispatchers.IO) {
|
||||
delay(2_100) // Looking busy...
|
||||
|
||||
logger.i("Navigation UI State is Landing")
|
||||
|
||||
Reference in New Issue
Block a user