Makes Client thread safe. Forcing the disconnect of an old relay list before connecting to a new one.

This commit is contained in:
Vitor Pamplona
2023-11-30 15:03:58 -05:00
parent eca420d789
commit 5db3ede09e
4 changed files with 48 additions and 24 deletions
@@ -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
}
@@ -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<HexKey>): Boolean {
@@ -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<Listener>()
private var relays = Constants.convertDefaultRelays()
private var relays = emptyArray<Relay>()
private var subscriptions = mapOf<String, List<TypedFilter>>()
@Synchronized
fun connect(relays: Array<Relay>) {
fun reconnect(relays: Array<Relay>?, 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<Relay>): Boolean {
if (relays.size != newRelayConfig.size) return false
fun isSameRelaySetConfig(newRelayConfig: Array<Relay>?): 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.
@@ -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()