feat(commons,desktop): wire latency tracker into RelayHealthStore + persistence

Phase 2 of relay-latency-health: hook the Phase 1 tracker into the existing
RelayHealthStore lifecycle and persist its rings via the existing
PreferencesRelayHealthPersistence so samples survive restarts.

commonMain:
  - RelayLatencyProvider: small interface so the store can drive a tracker
    that lives in jvmAndroidMain (the impl needs ConcurrentHashMap).
  - RelayHealthSnapshot: optional `latencySamples` field. Default empty;
    older saved snapshots load cleanly without it.
  - RelayHealthStore now takes optional `latencyTracker` / `nip11Provider`
    / `authProvider` constructor params:
      * exposes `latencySnapshots: StateFlow<ImmutableMap<Url, RelayLatencySnapshot>>`
        — MutableStateFlow updated inside the existing 60 s reclassify tick
        (one timer, not two — the tracker is scope-less and gets
        `sweep(now)` called from reclassify).
      * exposes `slowRelays: StateFlow<ImmutableMap<Url, SlowReason>>` —
        derived via `_latencySnapshots.map(classifySlowRelays).stateIn(
        scope, SharingStarted.Eagerly, persistentMapOf())`. The classifier
        reads `nip11Provider()` / `authProvider` live, so paid/auth-only
        relays only join the cohort once their auth completes.
      * `init {}` restores persisted samples into the tracker; the
        existing `schedulePersist()` now bundles `tracker.samplesForPersistence()`
        into the saved snapshot via a new private `snapshotForPersist()`
        helper. The same helper feeds the final flush in `close()`.
    No new dispatcher / scope / timer — everything piggybacks on the
    existing infra (single SupervisorJob, 5 s persist debounce, 60 s tick).

jvmAndroidMain:
  - RelayLatencyTracker now implements RelayLatencyProvider. Overrides drop
    the inline `System.currentTimeMillis()` default; callers from commonMain
    pass `TimeUtils.nowMillis()` explicitly.

desktopApp (jvmMain):
  - PreferencesRelayHealthPersistence persists per-relay latency rings in
    separate keys (`lat_<account-prefix>_<sha256(url)[..16]>`) so the 8 KB
    Preferences ceiling on the main `health_<account>` key isn't blown by a
    user with many relays. Each key holds one relay's four metric rings as
    `wss://relay.url\tok:csv|eose:csv|fr:csv|ping:csv`. On save, keys for
    relays no longer in the snapshot get removed so the prefs node doesn't
    grow unboundedly across account churn.

Notes:
  - Persistence still uses the existing 5 s debounce path. The deepened plan
    called for 30 s for `lat_*` keys; deferring that micro-optimization
    until we observe write thrash in practice. The cap on writes is
    one-rewrite-per-5s-of-activity which matches what the existing snooze
    persistence already does, so latency adds zero new flush events.
  - Tracker is wired only when a `RelayLatencyProvider` is passed to the
    store. Existing tests / Android continue to compile and run with
    latency unconfigured — `latencySnapshots` stays empty and `slowRelays`
    derives to empty. Desktop wiring lands in Phase 3.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
nrobi144
2026-06-18 12:11:25 +03:00
co-authored by Claude Opus 4.7
parent 96371715f4
commit a07c9ac7cc
5 changed files with 255 additions and 11 deletions
@@ -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<NormalizedRelayUrl, RelayHealthRecord> = emptyMap(),
val firstScanAt: Long = 0,
val lastSeenAny: Long = 0,
val latencySamples: Map<NormalizedRelayUrl, Map<LatencyMetric, IntArray>> = emptyMap(),
)
/**
@@ -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<NormalizedRelayUrl, Nip11RelayInformation?> = { 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<PersistentList<UnhealthyRelay>>(persistentListOf())
val unhealthy: StateFlow<PersistentList<UnhealthyRelay>> = _unhealthy.asStateFlow()
private val _latencySnapshots =
MutableStateFlow<ImmutableMap<NormalizedRelayUrl, RelayLatencySnapshot>>(persistentMapOf())
/** Per-relay rolling-window latency snapshot, updated on every 60 s tick. Empty when no tracker. */
val latencySnapshots: StateFlow<ImmutableMap<NormalizedRelayUrl, RelayLatencySnapshot>> =
_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<ImmutableMap<NormalizedRelayUrl, SlowReason>> =
_latencySnapshots
.map { snaps ->
if (latencyTracker == null || snaps.isEmpty()) {
persistentMapOf<NormalizedRelayUrl, SlowReason>()
} 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 {
@@ -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<NormalizedRelayUrl, RelayLatencySnapshot>
/** Raw per-relay per-metric sample arrays for persistence (chronological order). */
fun samplesForPersistence(): Map<NormalizedRelayUrl, Map<LatencyMetric, IntArray>>
/** Restore samples from a previous session — typically called once at app start. */
fun restoreSamples(saved: Map<NormalizedRelayUrl, Map<LatencyMetric, IntArray>>)
}
@@ -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<NormalizedRelayUrl, MutableMap<String, Long>>()
private val pendingSubId = ConcurrentHashMap<NormalizedRelayUrl, MutableMap<String, Long>>()
@@ -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<NormalizedRelayUrl, RelayLatencySnapshot> {
override fun snapshot(): ImmutableMap<NormalizedRelayUrl, RelayLatencySnapshot> {
if (samples.isEmpty()) return persistentMapOf()
val out = HashMap<NormalizedRelayUrl, RelayLatencySnapshot>(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_<url>` Preferences key.
*/
fun samplesForPersistence(): Map<NormalizedRelayUrl, Map<LatencyMetric, IntArray>> {
override fun samplesForPersistence(): Map<NormalizedRelayUrl, Map<LatencyMetric, IntArray>> {
if (samples.isEmpty()) return emptyMap()
val out = HashMap<NormalizedRelayUrl, Map<LatencyMetric, IntArray>>(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<NormalizedRelayUrl, Map<LatencyMetric, IntArray>>) {
override fun restoreSamples(saved: Map<NormalizedRelayUrl, Map<LatencyMetric, IntArray>>) {
for ((relay, perMetric) in saved) {
for ((metric, arr) in perMetric) {
if (arr.isEmpty()) continue
@@ -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_<account>_<hash>`) so the 8 KB
* Preferences ceiling on the main `health_<account>` 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<NormalizedRelayUrl, Map<LatencyMetric, IntArray>> {
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<NormalizedRelayUrl, Map<LatencyMetric, IntArray>>) {
val wantedKeys = mutableSetOf<String>()
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<LatencyMetric, IntArray>): 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<LatencyMetric, IntArray> {
if (raw.isBlank()) return emptyMap()
val out = mutableMapOf<LatencyMetric, IntArray>()
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()
}
}