Improve sync logic
This commit is contained in:
@@ -432,7 +432,9 @@ class DatabaseNostrRepository(
|
||||
storeNostrEventMutex.withLock {
|
||||
logger.d("saveNostrEvent: $nostrEvent")
|
||||
database.nostrDao().storeNostrEvent(
|
||||
nostrEvent,
|
||||
nostrEvent.copy(
|
||||
broadcastedAt = nostrEvent.createdAt
|
||||
),
|
||||
synchronizationRelayURLs = synchronizationRelayURLs,
|
||||
level = synchronizeNostrEventRequest.level
|
||||
)
|
||||
|
||||
@@ -75,6 +75,9 @@ interface NostrRepository {
|
||||
|
||||
suspend fun saveBroadcastReceipt(broadcastNostrEventReceipt: BroadcastNostrEventReceipt)
|
||||
|
||||
/**
|
||||
* Save nostrEvent from relays (and index for local viewing)...
|
||||
*/
|
||||
suspend fun saveNostrEvent(
|
||||
nostrEvent: NostrEvent,
|
||||
synchronizeNostrEventRequest: SynchronizeNostrEventRequest,
|
||||
|
||||
@@ -2,6 +2,7 @@ package ac.cord.auxiliary.compose.ui.view.model
|
||||
|
||||
import ac.cord.auxiliary.compose.database.model.BroadcastNostrEventRequest
|
||||
import ac.cord.auxiliary.compose.database.model.NostrEvent
|
||||
import ac.cord.auxiliary.compose.database.model.SynchronizeNostrEventRequest
|
||||
import ac.cord.auxiliary.compose.managers.SeedManager
|
||||
import ac.cord.auxiliary.compose.network.dto.RelayDTO
|
||||
import ac.cord.auxiliary.compose.network.relays.RelayPool.Companion.PUBLISH_TIMEOUT
|
||||
@@ -40,6 +41,7 @@ import kotlinx.coroutines.flow.distinctUntilChanged
|
||||
import kotlinx.coroutines.flow.getAndUpdate
|
||||
import kotlinx.coroutines.flow.timeout
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlin.time.Clock
|
||||
import kotlin.time.Duration.Companion.milliseconds
|
||||
import kotlin.time.Instant
|
||||
|
||||
@@ -303,19 +305,28 @@ class NavigationViewModel(
|
||||
)
|
||||
}
|
||||
is NostrIncomingMessage.NegentropyError -> {
|
||||
// TODO: Handle negentropy error...
|
||||
// nostrRepository.queueSynchronizeNostrEvent(
|
||||
// Relays.eventPublishRelaySet.take(1).map { normalizedRelay -> // TODO: Sync from all the publish relays...
|
||||
// SynchronizeNostrEventRequest(
|
||||
// purpose = "feed",
|
||||
// synchronizationFilters = arrayOf(
|
||||
// synchronizationFilter
|
||||
// ),
|
||||
// relayURL = normalizedRelay.url,
|
||||
// level = 0,
|
||||
// )
|
||||
// }
|
||||
// )
|
||||
logger.d("Negentropy Erroy (need to synchronize like normal): ${nostrIncomingMessage.negentropyReason}")
|
||||
|
||||
val since = negentropySynchronizeRequest.synchronizationFilter.since // TODO: Update since -> until based on what we have in the local db...
|
||||
val until = negentropySynchronizeRequest.synchronizationFilter.until
|
||||
|
||||
val synchronizationFilter = negentropySynchronizeRequest.synchronizationFilter.copy(
|
||||
since = since,
|
||||
until = until
|
||||
)
|
||||
|
||||
nostrRepository.queueSynchronizeNostrEvent(
|
||||
listOf(
|
||||
SynchronizeNostrEventRequest(
|
||||
purpose = negentropySynchronizeRequest.purpose,
|
||||
synchronizationFilters = arrayOf(
|
||||
synchronizationFilter
|
||||
),
|
||||
relayURL = negentropySynchronizeRequest.relayURL,
|
||||
level = negentropySynchronizeRequest.level,
|
||||
)
|
||||
)
|
||||
)
|
||||
}
|
||||
is NostrIncomingMessage.NegentropyMessage -> {
|
||||
logger.d("NegentropyMessage: ${nostrIncomingMessage.negentropyMessage}")
|
||||
@@ -326,6 +337,7 @@ class NavigationViewModel(
|
||||
logger.d("NeedIds: ${result.needIds.map { it.toHexString() }}")
|
||||
logger.d("SendIds: ${result.sendIds.map { it.toHexString() }}")
|
||||
logger.d("EventsIds: ${events.map { it.id }}")
|
||||
logger.d("Timestamp: ${events.map { it.createdAt.epochSeconds }}")
|
||||
|
||||
val broadcastNostrEventRequests = result.sendIds.map { nostrEventId ->
|
||||
BroadcastNostrEventRequest(
|
||||
@@ -365,30 +377,29 @@ class NavigationViewModel(
|
||||
localBroadcastNostrEventRequest.broadcastNostrEventRequest
|
||||
)
|
||||
|
||||
if (localBroadcastNostrEventRequest.nostrEvent.broadcastedAt == null) {
|
||||
scope.launch(Dispatchers.IO) {
|
||||
relaysSocketManager.publishEvent(
|
||||
localBroadcastNostrEventRequest.nostrEvent
|
||||
).timeout(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"
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
// if (localBroadcastNostrEventRequest.nostrEvent.broadcastedAt == null) {
|
||||
// scope.launch(Dispatchers.IO) {
|
||||
// relaysSocketManager.publishEvent(
|
||||
// localBroadcastNostrEventRequest.nostrEvent
|
||||
// ).timeout(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"
|
||||
// )
|
||||
// }
|
||||
// }
|
||||
// }
|
||||
// }
|
||||
|
||||
scope.launch(Dispatchers.IO) {
|
||||
// Broadcast to the intended relay...
|
||||
|
||||
Reference in New Issue
Block a user