From 57a64d258d7f56132166ef6cb944b5b62c6c5030 Mon Sep 17 00:00:00 2001 From: nrobi144 Date: Mon, 6 Jul 2026 16:08:18 +0300 Subject: [PATCH] fix(commons): guard RelayLatencyTracker.sweep against concurrent writes MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `putPending` stores each relay's pending map as `Collections.synchronizedMap(LinkedHashMap(...))` and wraps every writer path (`putPending`, `recordSent`, `recordIncoming`, `recordDisconnect`) in `synchronized(perRelay)`. `sweep` iterated `pending.entries.iterator()` without taking the same lock, violating the wrapper's Javadoc contract. Any concurrent websocket-thread write during `RelayHealthStore.reclassify`'s sweep threw `ConcurrentModificationException` on the underlying `LinkedHashMap$LinkedHashIterator`. Because `reclassify` schedules sweep on the AWT dispatcher, the CME killed Amethyst Desktop's Compose render thread and froze the UI. Wrap both inner iterator loops in `synchronized(pending) { ... }` — the exact synchronization the wrapper's Javadoc prescribes for manual iteration. Co-Authored-By: Claude Opus 4.7 (1M context) --- .../relays/health/RelayLatencyTracker.kt | 45 +++++++++++-------- 1 file changed, 26 insertions(+), 19 deletions(-) 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 43347efb20..44272f6535 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 @@ -166,31 +166,38 @@ class RelayLatencyTracker( * sample (per the brainstorm: "punish silent relays"). Idempotent and cheap. */ override fun sweep(nowMs: Long) { + // Per-relay pending maps are `Collections.synchronizedMap(LinkedHashMap)` — the + // Javadoc requires holding the wrapper's monitor while iterating, or a concurrent + // put/remove from putPending/recordSent/recordIncoming/recordDisconnect throws CME. for ((relay, pending) in pendingEventId) { - val it = pending.entries.iterator() - while (it.hasNext()) { - val (_, sentAt) = it.next() - if (nowMs - sentAt >= okTtlMs) { - ringFor(relay, LatencyMetric.OK_ACK).push(okTtlMs.toInt()) - it.remove() + synchronized(pending) { + val it = pending.entries.iterator() + while (it.hasNext()) { + val (_, sentAt) = it.next() + if (nowMs - sentAt >= okTtlMs) { + ringFor(relay, LatencyMetric.OK_ACK).push(okTtlMs.toInt()) + it.remove() + } } } } for ((relay, pending) in pendingSubId) { - val it = pending.entries.iterator() - while (it.hasNext()) { - val entry = it.next() - val sentAt = entry.value - if (nowMs - sentAt >= reqTtlMs) { - ringFor(relay, LatencyMetric.EOSE).push(reqTtlMs.toInt()) - // Record a FIRST_RESULT TTL sample only if the relay never emitted any - // matching event for this sub-id. - val seenSet = firstResultSeen[relay] - if (seenSet == null || entry.key !in seenSet) { - ringFor(relay, LatencyMetric.FIRST_RESULT).push(reqTtlMs.toInt()) + synchronized(pending) { + val it = pending.entries.iterator() + while (it.hasNext()) { + val entry = it.next() + val sentAt = entry.value + if (nowMs - sentAt >= reqTtlMs) { + ringFor(relay, LatencyMetric.EOSE).push(reqTtlMs.toInt()) + // Record a FIRST_RESULT TTL sample only if the relay never emitted any + // matching event for this sub-id. + val seenSet = firstResultSeen[relay] + if (seenSet == null || entry.key !in seenSet) { + ringFor(relay, LatencyMetric.FIRST_RESULT).push(reqTtlMs.toInt()) + } + seenSet?.remove(entry.key) + it.remove() } - seenSet?.remove(entry.key) - it.remove() } } }