Avoids race condition when update EOSEs

This commit is contained in:
Vitor Pamplona
2023-09-05 17:20:37 -04:00
parent 23414b0bc6
commit 6a3206b938
2 changed files with 13 additions and 1 deletions
@@ -14,6 +14,7 @@ import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel import kotlinx.coroutines.cancel
import kotlinx.coroutines.launch import kotlinx.coroutines.launch
import java.util.UUID import java.util.UUID
import java.util.concurrent.atomic.AtomicBoolean
import kotlin.Error import kotlin.Error
abstract class NostrDataSource(val debugName: String) { abstract class NostrDataSource(val debugName: String) {
@@ -23,6 +24,7 @@ abstract class NostrDataSource(val debugName: String) {
data class Counter(var counter: Int) data class Counter(var counter: Int)
private var eventCounter = mapOf<String, Counter>() private var eventCounter = mapOf<String, Counter>()
var changingFilters = AtomicBoolean()
fun printCounter() { fun printCounter() {
eventCounter.forEach { eventCounter.forEach {
@@ -59,7 +61,10 @@ abstract class NostrDataSource(val debugName: String) {
if (type == Relay.Type.EOSE && subscriptionId != null && subscriptionId in subscriptions.keys) { if (type == Relay.Type.EOSE && subscriptionId != null && subscriptionId in subscriptions.keys) {
// updates a per subscripton since date // updates a per subscripton since date
subscriptions[subscriptionId]?.updateEOSE(TimeUtils.now(), relay.url) subscriptions[subscriptionId]?.updateEOSE(
TimeUtils.fiveMinutesAgo(), // in case people's clock is slighly off.
relay.url
)
} }
} }
@@ -136,6 +141,8 @@ abstract class NostrDataSource(val debugName: String) {
// saves the current content to only update if it changes // saves the current content to only update if it changes
val currentFilters = activeSubscriptions.associate { it.id to it.toJson() } val currentFilters = activeSubscriptions.associate { it.id to it.toJson() }
changingFilters.getAndSet(true)
updateChannelFilters() updateChannelFilters()
// Makes sure to only send an updated filter when it actually changes. // Makes sure to only send an updated filter when it actually changes.
@@ -167,6 +174,8 @@ abstract class NostrDataSource(val debugName: String) {
} }
} }
} }
changingFilters.getAndSet(false)
} }
open fun consume(event: Event, relay: Relay) { open fun consume(event: Event, relay: Relay) {
@@ -134,6 +134,9 @@ object NostrSingleEventDataSource : NostrDataSource("SingleEventFeed") {
} }
val singleEventChannel = requestNewChannel { time, relayUrl -> val singleEventChannel = requestNewChannel { time, relayUrl ->
// Ignores EOSE if it is in the middle of a filter change.
if (changingFilters.get()) return@requestNewChannel
checkNotInMainThread() checkNotInMainThread()
eventsToWatch.forEach { eventsToWatch.forEach {