Bug fix on broadcast

This commit is contained in:
Kgothatso Ngako
2026-07-07 02:40:48 +02:00
parent 8368545f6f
commit ec919f5d9f
4 changed files with 32 additions and 30 deletions

View File

@@ -4,6 +4,7 @@ import androidx.room3.Dao
import androidx.room3.Insert
import androidx.room3.Query
import androidx.room3.Upsert
import at.torch.compose.database.model.intermdiate.LocalBroadcastNostrEventRequest
import kotlinx.coroutines.flow.Flow
import kotlin.time.Clock
import kotlin.time.Instant
@@ -14,7 +15,7 @@ interface BroadcastNostrEventRequestDao {
fun getAllBroadcastNostrEventRequests(): List<at.torch.compose.database.model.BroadcastNostrEventRequest>
@Query("SELECT * FROM BroadcastNostrEventRequest WHERE status = :status AND createdAt > :createdAt")
fun observeBroadcastNostrEventRequestsByStatus(status: String, createdAt: Instant = Clock.System.now()): Flow<List<at.torch.compose.database.model.intermdiate.LocalBroadcastNostrEventRequest>>
fun observeBroadcastNostrEventRequestsByStatus(status: String, createdAt: Instant = Clock.System.now()): Flow<LocalBroadcastNostrEventRequest?>
@Query("SELECT * FROM BroadcastNostrEventRequest WHERE nostrEventId = :nostrEventId")
fun observeBroadcastNostrEventRequestByNostrEventId(nostrEventId: String): Flow<at.torch.compose.database.model.BroadcastNostrEventRequest?>

View File

@@ -5,6 +5,7 @@ import at.torch.compose.database.model.BroadcastNostrEventRequest
import at.torch.compose.database.model.NostrEvent
import at.torch.compose.database.model.UnsignedNostrEvent
import at.torch.compose.database.model.intermdiate.LocalAccount
import at.torch.compose.database.model.intermdiate.LocalBroadcastNostrEventRequest
import at.torch.compose.database.model.types.SynchronizationFilter
import at.torch.compose.nostr.Relays
import co.touchlab.kermit.Logger
@@ -84,7 +85,7 @@ class DatabaseNostrRepository(
return database.unsignedNostrEventDao().observeUnsignedNostrEvents(publicKey)
}
override suspend fun observePendingBroadcastNostrEventRequests(): Flow<List<at.torch.compose.database.model.intermdiate.LocalBroadcastNostrEventRequest>> {
override suspend fun observePendingBroadcastNostrEventRequests(): Flow<LocalBroadcastNostrEventRequest?> {
return database.broadcastNostrEventRequestDao().observeBroadcastNostrEventRequestsByStatus("pending")
}

View File

@@ -32,7 +32,7 @@ interface NostrRepository {
suspend fun observeUnsignedNostrEvents(publicKey: HexKey): Flow<UnsignedNostrEvent?>
suspend fun observePendingBroadcastNostrEventRequests(): Flow<List<LocalBroadcastNostrEventRequest>>
suspend fun observePendingBroadcastNostrEventRequests(): Flow<LocalBroadcastNostrEventRequest?>
suspend fun observePendingSynchronizeNostrEventRequests(): Flow<SynchronizeNostrEventRequest?>
@@ -161,7 +161,7 @@ interface NostrRepository {
TODO("Not yet implemented")
}
override suspend fun observePendingBroadcastNostrEventRequests(): Flow<List<LocalBroadcastNostrEventRequest>> {
override suspend fun observePendingBroadcastNostrEventRequests(): Flow<LocalBroadcastNostrEventRequest> {
TODO("Not yet implemented")
}

View File

@@ -387,36 +387,37 @@ class SynchronizationViewModel(
logger.i { "observePendingBroadcastNostrEventRequests" }
scope.launch(Dispatchers.IO) {
nostrRepository.observePendingBroadcastNostrEventRequests().distinctUntilChanged().collect { localBroadcastNostrEventRequests ->
localBroadcastNostrEventRequests.forEach { localBroadcastNostrEventRequest ->
nostrRepository.observePendingBroadcastNostrEventRequests().distinctUntilChanged().collect { localBroadcastNostrEventRequest ->
logger.d("localBroadcastNostrEventRequest: $localBroadcastNostrEventRequest")
if (localBroadcastNostrEventRequest != null) {
nostrRepository.broadcastProcessed(
localBroadcastNostrEventRequest.broadcastNostrEventRequest
)
if (localBroadcastNostrEventRequest.nostrEvent.broadcastedAt == null) {
scope.launch(Dispatchers.IO) {
relaysSocketManager.publishEvent(
localBroadcastNostrEventRequest.nostrEvent
).timeout(RelayPool.PUBLISH_TIMEOUT.milliseconds).catch {
// Timeout...
}.collect { nostrPublishResult ->
if (nostrPublishResult.error != null) {
logger.e("Error publishing note: $nostrPublishResult")
nostrRepository.broadcastProcessed(
localBroadcastNostrEventRequest.broadcastNostrEventRequest,
"failed"
)
} else {
logger.d("Nostr Publish Result: $nostrPublishResult")
nostrRepository.broadcastProcessed(
localBroadcastNostrEventRequest.broadcastNostrEventRequest,
"published" // TODO: Might want to pass the nostrPublishResult.result
)
}
}
}
}
// if (localBroadcastNostrEventRequest.nostrEvent.broadcastedAt == null) {
// scope.launch(Dispatchers.IO) {
// relaysSocketManager.publishEvent(
// localBroadcastNostrEventRequest.nostrEvent
// ).timeout(RelayPool.PUBLISH_TIMEOUT.milliseconds).catch {
// // Timeout...
// }.collect { nostrPublishResult ->
// if (nostrPublishResult.error != null) {
// logger.e("Error publishing note: $nostrPublishResult")
// nostrRepository.broadcastProcessed(
// localBroadcastNostrEventRequest.broadcastNostrEventRequest,
// "failed"
// )
// } else {
// logger.d("Nostr Publish Result: $nostrPublishResult")
// nostrRepository.broadcastProcessed(
// localBroadcastNostrEventRequest.broadcastNostrEventRequest,
// "published" // TODO: Might want to pass the nostrPublishResult.result
// )
// }
// }
// }
// }
scope.launch(Dispatchers.IO) {
// Broadcast to the intended relay...
@@ -443,7 +444,6 @@ class SynchronizationViewModel(
}
}
}
}
}
}