mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
fix(commons): guard RelayLatencyTracker.sweep against concurrent writes
`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) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.7
parent
6361c54a4e
commit
57a64d258d
+26
-19
@@ -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()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user