mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 11:48:24 +00:00
feat(quartz): expose per-relay AUTH state as a Compose-stable StateFlow
RelayAuthStatus has to stay mutable — it holds LruCaches addressable from the per-relay OkHttp dispatcher thread, and replacing the whole holder on every mutation would be wasteful. But its mutability also makes it useless as a StateFlow value: mutating an entry doesn't change map identity, so distinct-until-changed downstream swallows the update and Compose never recomposes. Add an immutable view alongside: RelayAuthSnapshot (phase + lastAuthSuccessAt). RelayAuthStatus.snapshot() derives it from the LRU. RelayAuthenticator publishes a PersistentMap<NormalizedRelayUrl, RelayAuthSnapshot> via authStateFlow on every mutation (connect, disconnect, AUTH-submitted, AUTH-OK, AUTH-fail). PersistentMap gives O(log32 n) updates and a fresh identity per put, so both StateFlow equality and Compose strong-skipping work. This is the substrate for downstream consumers — the AUTH approval banner, the retry-queue wake on authCompleted, the indexer-fan-out gate — none of which are wired yet. They will read authStateFlow rather than querying RelayAuthStatus directly.
This commit is contained in:
+58
@@ -0,0 +1,58 @@
|
||||
/*
|
||||
* Copyright (c) 2025 Vitor Pamplona
|
||||
*
|
||||
* Permission is hereby granted, free of charge, to any person obtaining a copy of
|
||||
* this software and associated documentation files (the "Software"), to deal in
|
||||
* the Software without restriction, including without limitation the rights to use,
|
||||
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
|
||||
* Software, and to permit persons to whom the Software is furnished to do so,
|
||||
* subject to the following conditions:
|
||||
*
|
||||
* The above copyright notice and this permission notice shall be included in all
|
||||
* copies or substantial portions of the Software.
|
||||
*
|
||||
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
|
||||
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
|
||||
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
|
||||
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
|
||||
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
|
||||
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
||||
*/
|
||||
package com.vitorpamplona.quartz.nip01Core.relay.client.auth
|
||||
|
||||
import androidx.compose.runtime.Immutable
|
||||
|
||||
/**
|
||||
* Compose-stable per-relay AUTH snapshot exposed by [RelayAuthenticator].
|
||||
*
|
||||
* The internal [RelayAuthStatus] is a mutable holder around concurrent LRU
|
||||
* caches — necessary for the per-relay OkHttp dispatcher, but unsuitable as
|
||||
* a [kotlinx.coroutines.flow.StateFlow] value (mutating it doesn't change
|
||||
* identity, so distinct-until-changed swallows updates).
|
||||
*
|
||||
* [RelayAuthSnapshot] is the immutable view downstream consumers (UI banner,
|
||||
* retry coordinator, indexer-fan-out gate) subscribe to.
|
||||
*/
|
||||
@Immutable
|
||||
data class RelayAuthSnapshot(
|
||||
val phase: Phase,
|
||||
val lastAuthSuccessAt: Long?,
|
||||
) {
|
||||
enum class Phase {
|
||||
/** Connected; no AUTH challenge has been received yet. */
|
||||
IDLE,
|
||||
|
||||
/** Signed AUTH event in flight; awaiting OK from the relay. */
|
||||
AUTHENTICATING,
|
||||
|
||||
/** Last AUTH succeeded; relay accepts authenticated REQs. */
|
||||
AUTHENTICATED,
|
||||
|
||||
/** Last AUTH attempt failed; subsequent challenges may still arrive. */
|
||||
AUTH_FAILED,
|
||||
}
|
||||
|
||||
companion object {
|
||||
val IDLE = RelayAuthSnapshot(Phase.IDLE, lastAuthSuccessAt = null)
|
||||
}
|
||||
}
|
||||
+34
@@ -23,6 +23,8 @@ package com.vitorpamplona.quartz.nip01Core.relay.client.auth
|
||||
import androidx.collection.LruCache
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
import com.vitorpamplona.quartz.nip42RelayAuth.RelayAuthEvent
|
||||
import com.vitorpamplona.quartz.utils.TimeUtils
|
||||
import kotlin.concurrent.Volatile
|
||||
|
||||
class RelayAuthStatus {
|
||||
// Keeps track of auth responses to update the relay with all filters
|
||||
@@ -32,6 +34,12 @@ class RelayAuthStatus {
|
||||
// Avoids sending multiple replies for each auth.
|
||||
private val uniqueAuthChallengesSent: LruCache<ChallengePair, ChallengePair> = LruCache(10)
|
||||
|
||||
// Latest epoch-second at which a tracked AUTH event received a successful OK.
|
||||
// Read by RelayAuthSnapshot consumers for staleness checks (e.g. proactive
|
||||
// re-AUTH on window focus).
|
||||
@Volatile
|
||||
private var lastAuthSuccessAt: Long? = null
|
||||
|
||||
enum class AuthEventReceiptStatus {
|
||||
AUTHENTICATING,
|
||||
AUTHENTICATED,
|
||||
@@ -66,6 +74,7 @@ class RelayAuthStatus {
|
||||
return if (wasAlreadyAuthenticated != null) {
|
||||
if (success) {
|
||||
authResponseWatcher.put(eventId, AuthEventReceiptStatus.AUTHENTICATED)
|
||||
lastAuthSuccessAt = TimeUtils.now()
|
||||
} else {
|
||||
authResponseWatcher.put(eventId, AuthEventReceiptStatus.NOT_AUTHENTICATED)
|
||||
}
|
||||
@@ -77,4 +86,29 @@ class RelayAuthStatus {
|
||||
}
|
||||
|
||||
fun hasFinishedAllAuths() = authResponseWatcher.snapshot().all { it.value != AuthEventReceiptStatus.AUTHENTICATING }
|
||||
|
||||
/**
|
||||
* Build an immutable Compose-stable snapshot of the current per-relay AUTH
|
||||
* state. The phase is derived from the response watcher:
|
||||
*
|
||||
* - any AUTHENTICATING entry → [RelayAuthSnapshot.Phase.AUTHENTICATING]
|
||||
* - else any AUTHENTICATED entry → [RelayAuthSnapshot.Phase.AUTHENTICATED]
|
||||
* - else any NOT_AUTHENTICATED entry → [RelayAuthSnapshot.Phase.AUTH_FAILED]
|
||||
* - else (no tracked challenges) → [RelayAuthSnapshot.Phase.IDLE]
|
||||
*
|
||||
* The watcher LRU caps at 10 entries; a long-running connection that has
|
||||
* already AUTHed will still report AUTHENTICATED even after older entries
|
||||
* roll off, because the LRU keeps the most recent.
|
||||
*/
|
||||
fun snapshot(): RelayAuthSnapshot {
|
||||
val entries = authResponseWatcher.snapshot()
|
||||
val phase =
|
||||
when {
|
||||
entries.isEmpty() -> RelayAuthSnapshot.Phase.IDLE
|
||||
entries.values.any { it == AuthEventReceiptStatus.AUTHENTICATING } -> RelayAuthSnapshot.Phase.AUTHENTICATING
|
||||
entries.values.any { it == AuthEventReceiptStatus.AUTHENTICATED } -> RelayAuthSnapshot.Phase.AUTHENTICATED
|
||||
else -> RelayAuthSnapshot.Phase.AUTH_FAILED
|
||||
}
|
||||
return RelayAuthSnapshot(phase, lastAuthSuccessAt)
|
||||
}
|
||||
}
|
||||
|
||||
+38
-1
@@ -33,10 +33,16 @@ import com.vitorpamplona.quartz.nip01Core.signers.SignerExceptions
|
||||
import com.vitorpamplona.quartz.nip42RelayAuth.RelayAuthEvent
|
||||
import com.vitorpamplona.quartz.utils.Log
|
||||
import com.vitorpamplona.quartz.utils.cache.LargeCache
|
||||
import kotlinx.collections.immutable.PersistentMap
|
||||
import kotlinx.collections.immutable.persistentMapOf
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.IO
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.flow.MutableStateFlow
|
||||
import kotlinx.coroutines.flow.StateFlow
|
||||
import kotlinx.coroutines.flow.asStateFlow
|
||||
import kotlinx.coroutines.flow.update
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlin.coroutines.cancellation.CancellationException
|
||||
|
||||
@@ -61,8 +67,32 @@ class RelayAuthenticator(
|
||||
// Connection callbacks fire on the per-relay OkHttp dispatcher thread, so
|
||||
// this state is mutated concurrently — LargeCache wraps a platform-tuned
|
||||
// concurrent map (ConcurrentSkipListMap on jvmAndroid, CacheMap on Apple).
|
||||
//
|
||||
// This stays mutable because RelayAuthStatus carries an LruCache that has
|
||||
// to be addressable from the dispatcher thread. The Compose-observable
|
||||
// view of the same data is published on [authStateFlow] below, sourced
|
||||
// from RelayAuthStatus.snapshot().
|
||||
private val authStatus = LargeCache<NormalizedRelayUrl, RelayAuthStatus>()
|
||||
|
||||
private val _authStateFlow = MutableStateFlow<PersistentMap<NormalizedRelayUrl, RelayAuthSnapshot>>(persistentMapOf())
|
||||
|
||||
/**
|
||||
* Per-relay AUTH state as an immutable Compose-stable snapshot map.
|
||||
*
|
||||
* Downstream consumers (UI banner, retry queue, indexer-fan-out gate)
|
||||
* subscribe to this flow instead of polling [authStatus] directly.
|
||||
* Identity changes on every mutation, so [kotlinx.coroutines.flow.distinctUntilChanged]
|
||||
* downstream and Compose `@Immutable` skipping both work correctly.
|
||||
*/
|
||||
val authStateFlow: StateFlow<PersistentMap<NormalizedRelayUrl, RelayAuthSnapshot>> = _authStateFlow.asStateFlow()
|
||||
|
||||
private fun publishSnapshot(relayUrl: NormalizedRelayUrl) {
|
||||
val status = authStatus.get(relayUrl)
|
||||
_authStateFlow.update { current ->
|
||||
if (status == null) current.remove(relayUrl) else current.put(relayUrl, status.snapshot())
|
||||
}
|
||||
}
|
||||
|
||||
private val clientListener =
|
||||
object : RelayConnectionListener {
|
||||
override fun onIncomingMessage(
|
||||
@@ -78,10 +108,12 @@ class RelayAuthenticator(
|
||||
|
||||
override fun onConnecting(relay: IRelayClient) {
|
||||
authStatus.put(relay.url, RelayAuthStatus())
|
||||
publishSnapshot(relay.url)
|
||||
}
|
||||
|
||||
override fun onDisconnected(relay: IRelayClient) {
|
||||
authStatus.remove(relay.url)
|
||||
publishSnapshot(relay.url)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -102,6 +134,7 @@ class RelayAuthenticator(
|
||||
// only send replies to new challenges to avoid infinite loop:
|
||||
if (authStatus.get(relay.url)?.saveAuthSubmission(authEvent) == true) {
|
||||
relay.sendIfConnected(AuthCmd(authEvent))
|
||||
publishSnapshot(relay.url)
|
||||
}
|
||||
}
|
||||
} catch (e: CancellationException) {
|
||||
@@ -118,8 +151,12 @@ class RelayAuthenticator(
|
||||
relay: IRelayClient,
|
||||
msg: OkMessage,
|
||||
) {
|
||||
val transitioned = authStatus.get(relay.url)?.checkAuthResults(msg.eventId, msg.success) == true
|
||||
// Publish even on failure transitions so the UI can clear "AUTHENTICATING"
|
||||
// banners and reflect AUTH_FAILED state.
|
||||
publishSnapshot(relay.url)
|
||||
// if this is the OK of an auth event, renew all subscriptions and resend all outgoing events.
|
||||
if (authStatus.get(relay.url)?.checkAuthResults(msg.eventId, msg.success) == true) {
|
||||
if (transitioned) {
|
||||
client.syncFilters(relay)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user