diff --git a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/relays/health/RelayHealthPersistence.kt b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/relays/health/RelayHealthPersistence.kt index e16f6779f2..999026c7c1 100644 --- a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/relays/health/RelayHealthPersistence.kt +++ b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/relays/health/RelayHealthPersistence.kt @@ -27,11 +27,16 @@ import kotlin.concurrent.Volatile * Snapshot persisted to disk. Decoupled from the in-memory store so the same * persistence backend can swap (e.g. Android AccountSettings JSON vs Desktop * java.util.prefs.Preferences) without rippling into commons. + * + * [latencySamples] is populated only when a [RelayLatencyTracker] is wired into the store + * (Desktop today; Android follow-up). Persistence backends should treat the field as + * optional — older saved snapshots without it round-trip cleanly to empty maps. */ data class RelayHealthSnapshot( val records: Map = emptyMap(), val firstScanAt: Long = 0, val lastSeenAny: Long = 0, + val latencySamples: Map> = emptyMap(), ) /** diff --git a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/relays/health/RelayHealthStore.kt b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/relays/health/RelayHealthStore.kt index 7fba544322..7a6595a2f3 100644 --- a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/relays/health/RelayHealthStore.kt +++ b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/relays/health/RelayHealthStore.kt @@ -21,9 +21,12 @@ package com.vitorpamplona.amethyst.commons.relays.health import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl +import com.vitorpamplona.quartz.nip11RelayInfo.Nip11RelayInformation import com.vitorpamplona.quartz.utils.TimeUtils +import kotlinx.collections.immutable.ImmutableMap import kotlinx.collections.immutable.PersistentList import kotlinx.collections.immutable.persistentListOf +import kotlinx.collections.immutable.persistentMapOf import kotlinx.coroutines.CoroutineDispatcher import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers @@ -32,8 +35,11 @@ import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.cancel import kotlinx.coroutines.delay import kotlinx.coroutines.flow.MutableStateFlow +import kotlinx.coroutines.flow.SharingStarted import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.asStateFlow +import kotlinx.coroutines.flow.map +import kotlinx.coroutines.flow.stateIn import kotlinx.coroutines.flow.update import kotlinx.coroutines.launch import kotlinx.coroutines.withContext @@ -61,6 +67,16 @@ class RelayHealthStore( // (Android/Desktop) should pass `Dispatchers.IO` so `prefs.flush()` doesn't sit on // a CPU-bound worker. private val ioDispatcher: CoroutineDispatcher = Dispatchers.Default, + // Optional latency tracker (jvmAndroidMain in the current build — Desktop wires it, + // Android picks up later). When null, [latencySnapshots] stays empty and [slowRelays] + // never flags anything. + private val latencyTracker: RelayLatencyProvider? = null, + // NIP-11 info supplier, used by the slow-relay classifier to skip auth/payment-required + // relays. Read on every reclassify tick, so the supplier can return a live snapshot. + private val nip11Provider: () -> ImmutableMap = { persistentMapOf() }, + // Per-relay auth-completed lookup. Used by the classifier alongside [nip11Provider] — + // a NIP-11 auth-required relay only participates in the cohort once auth is complete. + private val authProvider: (NormalizedRelayUrl) -> Boolean = { true }, ) { companion object { const val PERSIST_DEBOUNCE_MS: Long = 5_000L @@ -89,6 +105,35 @@ class RelayHealthStore( private val _unhealthy = MutableStateFlow>(persistentListOf()) val unhealthy: StateFlow> = _unhealthy.asStateFlow() + private val _latencySnapshots = + MutableStateFlow>(persistentMapOf()) + + /** Per-relay rolling-window latency snapshot, updated on every 60 s tick. Empty when no tracker. */ + val latencySnapshots: StateFlow> = + _latencySnapshots.asStateFlow() + + /** + * Relays classified as "slow" — derived from [latencySnapshots] + the NIP-11 / auth providers. + * Re-runs whenever [latencySnapshots] emits (every 60 s tick when a tracker is wired). + * + * `SharingStarted.Eagerly` so the value is fresh whenever the UI reads `.value` without + * an active collector — matches how callers consume [unhealthy]. + */ + val slowRelays: StateFlow> = + _latencySnapshots + .map { snaps -> + if (latencyTracker == null || snaps.isEmpty()) { + persistentMapOf() + } else { + classifySlowRelays( + snapshots = snaps, + nip11 = nip11Provider(), + torEnabled = torEnabledProvider(), + authStatus = authProvider, + ) + } + }.stateIn(scope, SharingStarted.Eagerly, persistentMapOf()) + @Volatile private var persistJob: Job? = null @Volatile private var tickJob: Job? = null @@ -96,6 +141,15 @@ class RelayHealthStore( @Volatile private var closed = false init { + // Restore persisted latency samples into the tracker (if both wired). + if (latencyTracker != null) { + val loadedLatency = state.value.latencySamples + if (loadedLatency.isNotEmpty()) { + latencyTracker.restoreSamples(loadedLatency) + _latencySnapshots.value = latencyTracker.snapshot() + } + } + // Persist the firstScanAt seed if we just stamped it. schedulePersist() @@ -209,6 +263,17 @@ class RelayHealthStore( val flagged = withContext(ioDispatcher) { + // Sweep + snapshot the latency tracker inside the same tick (one timer, not two). + if (latencyTracker != null) { + latencyTracker.sweep(TimeUtils.nowMillis()) + val snap = latencyTracker.snapshot() + // Only emit if anything changed structurally — keeps strong-skipping happy + // on the dashboard rows when the per-tick snapshot is value-equal to the + // previous one (e.g. quiet period, no new samples). + if (snap != _latencySnapshots.value) { + _latencySnapshots.value = snap + } + } classifyRelayHealth( records = s.records, listMembership = membership, @@ -227,13 +292,23 @@ class RelayHealthStore( persistJob = scope.launch { delay(PERSIST_DEBOUNCE_MS) - val snapshot = state.value withContext(ioDispatcher) { - runCatching { persistence.save(snapshot) } + runCatching { persistence.save(snapshotForPersist()) } } } } + /** + * Bundles the timestamp/snooze state with the tracker's current sample arrays. Read at + * save time so we don't have to mirror per-event sample updates into `state` (which + * would break StateFlow equality with IntArrays). + */ + private fun snapshotForPersist(): RelayHealthSnapshot { + val timestamps = state.value + val latency = latencyTracker?.samplesForPersistence().orEmpty() + return if (latency.isEmpty()) timestamps else timestamps.copy(latencySamples = latency) + } + /** * Tear down internal jobs and fire the final persist off-thread. Safe to call from * the composition / Main thread: the blocking I/O is dispatched to [ioDispatcher] @@ -245,7 +320,7 @@ class RelayHealthStore( closed = true persistJob?.cancel() tickJob?.cancel() - val finalSnapshot = state.value + val finalSnapshot = snapshotForPersist() val flushScope = CoroutineScope(SupervisorJob() + ioDispatcher) flushScope.launch { try { diff --git a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/relays/health/RelayLatencyProvider.kt b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/relays/health/RelayLatencyProvider.kt new file mode 100644 index 0000000000..6281a8adf4 --- /dev/null +++ b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/relays/health/RelayLatencyProvider.kt @@ -0,0 +1,46 @@ +/* + * 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.amethyst.commons.relays.health + +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl +import com.vitorpamplona.quartz.utils.TimeUtils +import kotlinx.collections.immutable.ImmutableMap + +/** + * commonMain-facing view of a relay latency tracker. The [RelayHealthStore] holds a + * reference to this interface and drives `sweep` + `snapshot` from its existing 60 s + * reclassify tick. The concrete implementation lives in `jvmAndroidMain` because it + * needs `ConcurrentHashMap` for pending-correlation maps; iOS picks up later via an + * `expect class` once that target is wired. + */ +interface RelayLatencyProvider { + /** Expire pending entries past their TTL, recording the TTL value as the sample. */ + fun sweep(nowMs: Long = TimeUtils.nowMillis()) + + /** Immutable snapshot of all tracked relays' rolling-window medians. */ + fun snapshot(): ImmutableMap + + /** Raw per-relay per-metric sample arrays for persistence (chronological order). */ + fun samplesForPersistence(): Map> + + /** Restore samples from a previous session — typically called once at app start. */ + fun restoreSamples(saved: Map>) +} diff --git a/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/relays/health/RelayLatencyTracker.kt b/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/relays/health/RelayLatencyTracker.kt index 07e3714416..43347efb20 100644 --- a/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/relays/health/RelayLatencyTracker.kt +++ b/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/relays/health/RelayLatencyTracker.kt @@ -70,7 +70,7 @@ class RelayLatencyTracker( val maxPendingPerRelay: Int = DEFAULT_MAX_PENDING_PER_RELAY, val okTtlMs: Long = DEFAULT_OK_TTL_MS, val reqTtlMs: Long = DEFAULT_REQ_TTL_MS, -) { +) : RelayLatencyProvider { // Per-relay pending maps. `Long` is `currentTimeMillis()` at the time of the send. private val pendingEventId = ConcurrentHashMap>() private val pendingSubId = ConcurrentHashMap>() @@ -165,7 +165,7 @@ class RelayLatencyTracker( * Expires pending entries older than the configured TTLs and records the TTL value as the * sample (per the brainstorm: "punish silent relays"). Idempotent and cheap. */ - fun sweep(nowMs: Long = System.currentTimeMillis()) { + override fun sweep(nowMs: Long) { for ((relay, pending) in pendingEventId) { val it = pending.entries.iterator() while (it.hasNext()) { @@ -199,7 +199,7 @@ class RelayLatencyTracker( // ------ Snapshot ------ /** Immutable snapshot of all tracked relays' current rolling-window medians. */ - fun snapshot(): ImmutableMap { + override fun snapshot(): ImmutableMap { if (samples.isEmpty()) return persistentMapOf() val out = HashMap(samples.size) for ((relay, perMetric) in samples) { @@ -219,7 +219,7 @@ class RelayLatencyTracker( * Raw per-relay per-metric sample arrays (chronological order). Used by the persistence * layer to encode the rings into the packed `lat_` Preferences key. */ - fun samplesForPersistence(): Map> { + override fun samplesForPersistence(): Map> { if (samples.isEmpty()) return emptyMap() val out = HashMap>(samples.size) for ((relay, perMetric) in samples) { @@ -238,7 +238,7 @@ class RelayLatencyTracker( * existing buffers for the listed `(relay, metric)` pairs; other relays/metrics are * untouched. */ - fun restoreSamples(saved: Map>) { + override fun restoreSamples(saved: Map>) { for ((relay, perMetric) in saved) { for ((metric, arr) in perMetric) { if (arr.isEmpty()) continue diff --git a/desktopApp/src/jvmMain/kotlin/com/vitorpamplona/amethyst/desktop/model/PreferencesRelayHealthPersistence.kt b/desktopApp/src/jvmMain/kotlin/com/vitorpamplona/amethyst/desktop/model/PreferencesRelayHealthPersistence.kt index e04301ff69..8e7ce912ca 100644 --- a/desktopApp/src/jvmMain/kotlin/com/vitorpamplona/amethyst/desktop/model/PreferencesRelayHealthPersistence.kt +++ b/desktopApp/src/jvmMain/kotlin/com/vitorpamplona/amethyst/desktop/model/PreferencesRelayHealthPersistence.kt @@ -20,11 +20,14 @@ */ package com.vitorpamplona.amethyst.desktop.model +import com.vitorpamplona.amethyst.commons.relays.health.LatencyMetric import com.vitorpamplona.amethyst.commons.relays.health.RelayHealthPersistence import com.vitorpamplona.amethyst.commons.relays.health.RelayHealthRecord import com.vitorpamplona.amethyst.commons.relays.health.RelayHealthSnapshot import com.vitorpamplona.quartz.nip01Core.core.HexKey +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer +import java.security.MessageDigest import java.util.prefs.Preferences /** @@ -50,11 +53,13 @@ class PreferencesRelayHealthPersistence( private val storageKey: String = "health_${userPubKeyHex.take(8)}" + private val latencyKeyPrefix: String = "lat_${userPubKeyHex.take(8)}_" + override fun load(): RelayHealthSnapshot { return try { - val raw = prefs.get(storageKey, null) ?: return RelayHealthSnapshot() + val raw = prefs.get(storageKey, null) ?: return RelayHealthSnapshot(latencySamples = loadLatencySamples()) val lines = raw.split('\n') - if (lines.size < 2) return RelayHealthSnapshot() + if (lines.size < 2) return RelayHealthSnapshot(latencySamples = loadLatencySamples()) val firstScanAt = lines[0].toLongOrNull() ?: 0L val lastSeenAny = lines[1].toLongOrNull() ?: 0L val records = @@ -72,7 +77,7 @@ class PreferencesRelayHealthPersistence( put(url, rec) } } - RelayHealthSnapshot(records, firstScanAt, lastSeenAny) + RelayHealthSnapshot(records, firstScanAt, lastSeenAny, loadLatencySamples()) } catch (_: Exception) { RelayHealthSnapshot() } @@ -92,12 +97,125 @@ class PreferencesRelayHealthPersistence( if (sb.length > MAX_PREFS_VALUE_LENGTH) break } prefs.put(storageKey, sb.toString().take(MAX_PREFS_VALUE_LENGTH)) + saveLatencySamples(snapshot.latencySamples) prefs.flush() } catch (_: Exception) { } } + // ------ latency samples ------ + + /** + * Per-relay latency rings live in their own keys (`lat__`) so the 8 KB + * Preferences ceiling on the main `health_` key isn't blown by a heavy user. + * Each key holds one relay's four metric rings: `wss://relay.url\tok:csv|eose:csv|fr:csv|ping:csv`. + */ + private fun loadLatencySamples(): Map> { + return try { + val keys = prefs.keys() + if (keys.isEmpty()) return emptyMap() + buildMap { + for (k in keys) { + if (!k.startsWith(latencyKeyPrefix)) continue + val raw = prefs.get(k, null) ?: continue + val tab = raw.indexOf('\t') + if (tab <= 0) continue + val url = RelayUrlNormalizer.normalizeOrNull(raw.substring(0, tab)) ?: continue + val perMetric = decodeMetricCsv(raw.substring(tab + 1)) + if (perMetric.isNotEmpty()) put(url, perMetric) + } + } + } catch (_: Exception) { + emptyMap() + } + } + + private fun saveLatencySamples(samples: Map>) { + val wantedKeys = mutableSetOf() + for ((url, perMetric) in samples) { + if (perMetric.isEmpty()) continue + val key = latencyKeyPrefix + relayKeySuffix(url) + wantedKeys.add(key) + val packed = url.url + "\t" + encodeMetricCsv(perMetric) + // 8 KB Preferences ceiling per key; 50 samples × 4 metrics × ~5 chars ≈ 1 KB, + // so single relays never come close. Truncate as a defensive measure only. + prefs.put(key, packed.take(MAX_PREFS_VALUE_LENGTH)) + } + // Drop keys for relays that are no longer in the snapshot (relay removed from account). + for (k in prefs.keys()) { + if (k.startsWith(latencyKeyPrefix) && k !in wantedKeys) prefs.remove(k) + } + } + + private fun encodeMetricCsv(perMetric: Map): String { + val sb = StringBuilder() + var first = true + for (metric in LatencyMetric.entries) { + val arr = perMetric[metric] ?: continue + if (arr.isEmpty()) continue + if (!first) sb.append('|') + first = false + sb.append(metricTag(metric)).append(':') + for ((i, v) in arr.withIndex()) { + if (i > 0) sb.append(',') + sb.append(v) + } + } + return sb.toString() + } + + private fun decodeMetricCsv(raw: String): Map { + if (raw.isBlank()) return emptyMap() + val out = mutableMapOf() + for (section in raw.split('|')) { + val colon = section.indexOf(':') + if (colon <= 0) continue + val metric = parseMetricTag(section.substring(0, colon)) ?: continue + val csv = section.substring(colon + 1) + if (csv.isBlank()) continue + val parts = csv.split(',') + val arr = IntArray(parts.size) + var n = 0 + for (p in parts) { + val v = p.toIntOrNull() ?: continue + arr[n++] = v + } + out[metric] = if (n == arr.size) arr else arr.copyOf(n) + } + return out + } + + private fun metricTag(metric: LatencyMetric): String = + when (metric) { + LatencyMetric.OK_ACK -> "ok" + LatencyMetric.EOSE -> "eose" + LatencyMetric.FIRST_RESULT -> "fr" + LatencyMetric.PING -> "ping" + } + + private fun parseMetricTag(tag: String): LatencyMetric? = + when (tag) { + "ok" -> LatencyMetric.OK_ACK + "eose" -> LatencyMetric.EOSE + "fr" -> LatencyMetric.FIRST_RESULT + "ping" -> LatencyMetric.PING + else -> null + } + + /** Short stable suffix from the relay URL. SHA-256 first 16 hex chars — 64 bits of entropy. */ + private fun relayKeySuffix(url: NormalizedRelayUrl): String { + val digest = MessageDigest.getInstance("SHA-256").digest(url.url.toByteArray(Charsets.UTF_8)) + val sb = StringBuilder(16) + for (i in 0 until 8) { + val b = digest[i].toInt() and 0xff + sb.append(HEX[b shr 4]) + sb.append(HEX[b and 0xf]) + } + return sb.toString() + } + companion object { private const val MAX_PREFS_VALUE_LENGTH = 8000 + private val HEX = "0123456789abcdef".toCharArray() } }