From 5db3ede09e8a17899ed61c36e106e7911fc944bd Mon Sep 17 00:00:00 2001 From: Vitor Pamplona Date: Thu, 30 Nov 2023 15:03:58 -0500 Subject: [PATCH] Makes Client thread safe. Forcing the disconnect of an old relay list before connecting to a new one. --- .../vitorpamplona/amethyst/ServiceManager.kt | 6 ++- .../vitorpamplona/amethyst/model/Account.kt | 7 +-- .../amethyst/service/relays/Client.kt | 49 +++++++++++++------ .../amethyst/service/relays/Relay.kt | 10 +++- 4 files changed, 48 insertions(+), 24 deletions(-) diff --git a/app/src/main/java/com/vitorpamplona/amethyst/ServiceManager.kt b/app/src/main/java/com/vitorpamplona/amethyst/ServiceManager.kt index 36907086c..53f01f4b4 100644 --- a/app/src/main/java/com/vitorpamplona/amethyst/ServiceManager.kt +++ b/app/src/main/java/com/vitorpamplona/amethyst/ServiceManager.kt @@ -72,7 +72,9 @@ class ServiceManager { } if (myAccount != null) { - Client.connect(myAccount.activeRelays() ?: myAccount.convertLocalRelays()) + val relaySet = myAccount.activeRelays() ?: myAccount.convertLocalRelays() + Log.d("Relay", "Service Manager Connect Connecting ${relaySet.size}") + Client.reconnect(relaySet) // start services NostrAccountDataSource.account = myAccount @@ -120,7 +122,7 @@ class ServiceManager { NostrUserProfileDataSource.stopSync() NostrVideoDataSource.stopSync() - Client.disconnect() + Client.reconnect(null) isStarted = false } diff --git a/app/src/main/java/com/vitorpamplona/amethyst/model/Account.kt b/app/src/main/java/com/vitorpamplona/amethyst/model/Account.kt index 404f85b49..59cccf609 100644 --- a/app/src/main/java/com/vitorpamplona/amethyst/model/Account.kt +++ b/app/src/main/java/com/vitorpamplona/amethyst/model/Account.kt @@ -19,7 +19,6 @@ import com.vitorpamplona.amethyst.service.relays.Client import com.vitorpamplona.amethyst.service.relays.Constants import com.vitorpamplona.amethyst.service.relays.FeedType import com.vitorpamplona.amethyst.service.relays.Relay -import com.vitorpamplona.amethyst.service.relays.RelayPool import com.vitorpamplona.amethyst.ui.components.BundledUpdate import com.vitorpamplona.quartz.crypto.KeyPair import com.vitorpamplona.quartz.encoders.ATag @@ -1809,11 +1808,7 @@ class Account( fun reconnectIfRelaysHaveChanged() { val newRelaySet = activeRelays() ?: convertLocalRelays() - if (!Client.isSameRelaySetConfig(newRelaySet)) { - Client.disconnect() - Client.connect(newRelaySet) - RelayPool.requestAndWatch() - } + Client.reconnect(newRelaySet, true) } fun isAllHidden(users: Set): Boolean { diff --git a/app/src/main/java/com/vitorpamplona/amethyst/service/relays/Client.kt b/app/src/main/java/com/vitorpamplona/amethyst/service/relays/Client.kt index 3c2047c0c..156ea4d3a 100644 --- a/app/src/main/java/com/vitorpamplona/amethyst/service/relays/Client.kt +++ b/app/src/main/java/com/vitorpamplona/amethyst/service/relays/Client.kt @@ -1,5 +1,6 @@ package com.vitorpamplona.amethyst.service.relays +import android.util.Log import com.vitorpamplona.amethyst.service.HttpClient import com.vitorpamplona.amethyst.service.checkNotInMainThread import com.vitorpamplona.quartz.events.Event @@ -18,21 +19,47 @@ import java.util.UUID */ object Client : RelayPool.Listener { private var listeners = setOf() - private var relays = Constants.convertDefaultRelays() + private var relays = emptyArray() private var subscriptions = mapOf>() @Synchronized - fun connect(relays: Array) { + fun reconnect(relays: Array?, onlyIfChanged: Boolean = false) { + Log.d("Relay", "Relay Pool Reconnecting to ${relays?.size} relays") checkNotInMainThread() - RelayPool.register(this) - RelayPool.unloadRelays() - RelayPool.loadRelays(relays.toList()) - this.relays = relays + if (onlyIfChanged) { + if (!isSameRelaySetConfig(relays)) { + if (this.relays.isNotEmpty()) { + RelayPool.disconnect() + RelayPool.unregister(this) + RelayPool.unloadRelays() + } + + if (relays != null) { + RelayPool.register(this) + RelayPool.loadRelays(relays.toList()) + RelayPool.requestAndWatch() + this.relays = relays + } + } + } else { + if (this.relays.isNotEmpty()) { + RelayPool.disconnect() + RelayPool.unregister(this) + RelayPool.unloadRelays() + } + + if (relays != null) { + RelayPool.register(this) + RelayPool.loadRelays(relays.toList()) + RelayPool.requestAndWatch() + this.relays = relays + } + } } - fun isSameRelaySetConfig(newRelayConfig: Array): Boolean { - if (relays.size != newRelayConfig.size) return false + fun isSameRelaySetConfig(newRelayConfig: Array?): Boolean { + if (relays.size != newRelayConfig?.size) return false relays.forEach { oldRelayInfo -> val newRelayInfo = newRelayConfig.find { it.url == oldRelayInfo.url } ?: return false @@ -126,12 +153,6 @@ object Client : RelayPool.Listener { subscriptions = subscriptions.minus(subscriptionId) } - fun disconnect() { - RelayPool.unregister(this) - RelayPool.disconnect() - RelayPool.unloadRelays() - } - @OptIn(DelicateCoroutinesApi::class) override fun onEvent(event: Event, subscriptionId: String, relay: Relay) { // Releases the Web thread for the new payload. diff --git a/app/src/main/java/com/vitorpamplona/amethyst/service/relays/Relay.kt b/app/src/main/java/com/vitorpamplona/amethyst/service/relays/Relay.kt index 2344a2dd8..2db0bd446 100644 --- a/app/src/main/java/com/vitorpamplona/amethyst/service/relays/Relay.kt +++ b/app/src/main/java/com/vitorpamplona/amethyst/service/relays/Relay.kt @@ -87,7 +87,7 @@ class Relay( private var connectingBlock = AtomicBoolean() fun connectAndRun(onConnected: (Relay) -> Unit) { - Log.d("Client", "Relay.connect $url") + Log.d("Relay", "Relay.connect $url") // BRB is crashing OkHttp Deflater object :( if (url.contains("brb.io")) return @@ -120,6 +120,7 @@ class Relay( inner class RelayListener(val onConnected: (Relay) -> Unit) : WebSocketListener() { override fun onOpen(webSocket: WebSocket, response: Response) { checkNotInMainThread() + Log.d("Relay", "Connect onOpen $url $socket") markConnectionAsReady( pingInMs = response.receivedResponseAtMillis - response.sentRequestAtMillis, @@ -150,6 +151,8 @@ class Relay( override fun onClosing(webSocket: WebSocket, code: Int, reason: String) { checkNotInMainThread() + Log.w("Relay", "Relay onClosing $url: $reason") + listeners.forEach { it.onRelayStateChange( this@Relay, @@ -164,6 +167,8 @@ class Relay( markConnectionAsClosed() + Log.w("Relay", "Relay onClosed $url: $reason") + listeners.forEach { it.onRelayStateChange(this@Relay, StateType.DISCONNECT, null) } } @@ -211,6 +216,7 @@ class Relay( listeners.forEach { it.onEvent(this@Relay, subscriptionId, event) if (afterEOSEPerSubscription[subscriptionId] == true) { + Log.w("Relay", "Relay onEOSE $url $subscriptionId") it.onRelayStateChange(this@Relay, StateType.EOSE, subscriptionId) } } @@ -267,7 +273,7 @@ class Relay( } fun disconnect() { - Log.d("Client", "Relay.disconnect $url") + Log.d("Relay", "Relay.disconnect $url") checkNotInMainThread() closingTimeInSeconds = TimeUtils.now()