mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 03:38:23 +00:00
fix: age the dedup cache on a clock, not on idleness
The idle-triggered release in the previous commit never fired. Measured on the emulator: 240s of a completely flat heap, and every tick returned false. Cause: relays keep pushing events down open subscriptions forever, so the decoder is never idle in the sense of "saw no frames". Counting cache hits as activity -- which the tests asserted, and which is right in isolation, since a stream of pure duplicates is when the cache is earning its keep -- guaranteed the idle clock was refreshed indefinitely. The unit tests proved the intended behaviour; only the device showed the intent was wrong. Replaced with unconditional clock-based aging: ageOutCache() retires the live generation, so an id survives one tick and dies on the next, and a 30s tick bounds any id's lifetime at ~60s. clearCache() still releases everything when the host disconnects. Capacity keeps bounding cost during a burst; this bounds how long it costs anything afterwards. On-device A/B, same account and duration: idle-based (never fired): flat 127-128 MB for 240s clock-based, run 1: 130 MB -> 71 MB when 7,436 ids were dropped clock-based, run 2: 143 MB -> 121 MB (3,153 ids) -> 96 MB (4,530 ids) Each run contains its own control: the first tick only retires a generation, drops nothing, and frees nothing. Note ~7-8 KB released per cached id, an order above the ~600 B/entry the recorded corpus suggested -- live traffic carries big kind-0/kind-3 events the capture's first 150k frames under-represent. So capacity 8192 was costing tens of MB on a real account, not the 4.7 MB the corpus implied. Tests rewritten for the new semantics (no clock needed now, so they are fully deterministic) and mutation-checked: a no-op ageOut fails 4, an ageOut that clears both generations at once fails 3. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5
parent
1b8e264706
commit
ccebc3d344
+16
-13
@@ -160,7 +160,7 @@ class NostrClient(
|
||||
// decoder cached has no dedup value left across that gap, but it would keep
|
||||
// costing memory for the whole background stretch, so release it now rather
|
||||
// than leaving a timer running to do it later.
|
||||
decoder.trimIfIdle(idleMillis = 0)
|
||||
decoder.clearCache()
|
||||
}
|
||||
|
||||
override fun isActive() = isActiveFlow.value
|
||||
@@ -214,30 +214,33 @@ class NostrClient(
|
||||
}
|
||||
|
||||
/**
|
||||
* Releases the decoder's cache after a quiet stretch.
|
||||
* Ages out the decoder's cache on a clock.
|
||||
*
|
||||
* A caching decoder bounds what it holds by capacity, which caps the cost during
|
||||
* a burst but never ends it: its generations only rotate inside an insert, so a
|
||||
* client that goes quiet pins every cached event indefinitely. Dedup value decays
|
||||
* within seconds (duplicates are the same event arriving from other relays), so
|
||||
* anything still held after [DECODER_IDLE_MS] is pure cost.
|
||||
* A caching decoder bounds what it holds by capacity, which caps the cost during a
|
||||
* burst but never ends it: its generations only rotate inside an insert, so a
|
||||
* client that drops to a trickle keeps everything from the initial burst forever.
|
||||
*
|
||||
* Suspends on [isActiveFlow] like [keepAliveJob], so no timer fires while the
|
||||
* client is down — [disconnect] already trims eagerly on the way there.
|
||||
* Deliberately NOT idle-triggered. Measured on device: relays keep pushing events
|
||||
* down open subscriptions indefinitely, so the decoder is never idle in the sense
|
||||
* of "saw no frames" — an idle check returned false on every tick across 240s of
|
||||
* a flat heap. Aging unconditionally is what actually releases the memory.
|
||||
*
|
||||
* Ids survive one tick and die on the next, so [DECODER_AGE_OUT_MS] bounds their
|
||||
* lifetime at ~2x that. Suspends on [isActiveFlow] like [keepAliveJob], so no timer
|
||||
* fires while the client is down — [disconnect] clears outright on the way there.
|
||||
*/
|
||||
private val decoderTrimJob =
|
||||
scope.launch {
|
||||
while (true) {
|
||||
isActiveFlow.first { it }
|
||||
delay(DECODER_TRIM_INTERVAL_MS)
|
||||
decoder.trimIfIdle(DECODER_IDLE_MS)
|
||||
delay(DECODER_AGE_OUT_MS)
|
||||
decoder.ageOutCache()
|
||||
}
|
||||
}
|
||||
|
||||
companion object {
|
||||
private const val KEEP_ALIVE_INTERVAL_MS = 60_000L
|
||||
private const val DECODER_TRIM_INTERVAL_MS = 30_000L
|
||||
private const val DECODER_IDLE_MS = 60_000L
|
||||
private const val DECODER_AGE_OUT_MS = 30_000L
|
||||
}
|
||||
|
||||
override fun reconnect(
|
||||
|
||||
+22
-29
@@ -22,7 +22,6 @@ package com.vitorpamplona.quartz.nip01Core.relay.commands.toClient
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
import com.vitorpamplona.quartz.utils.TimeUtils
|
||||
import com.vitorpamplona.quartz.utils.cache.ConcurrentHashCache
|
||||
import kotlin.concurrent.Volatile
|
||||
import kotlin.concurrent.atomics.AtomicLong
|
||||
@@ -94,16 +93,6 @@ class CachingEventDecoder(
|
||||
|
||||
@Volatile private var previous = ConcurrentHashCache<HexKey, Event>()
|
||||
|
||||
/**
|
||||
* When an EVENT frame was last seen, for [trimIfIdle].
|
||||
*
|
||||
* Written on every EVENT frame — hit or miss — so a cache that is actively
|
||||
* serving duplicates counts as busy and is never trimmed out from under the
|
||||
* traffic it is helping. One relaxed volatile store per frame; at observed
|
||||
* rates (~400 frames/s) that is not measurable.
|
||||
*/
|
||||
@Volatile private var lastActivityAt: Long = TimeUtils.nowMillis()
|
||||
|
||||
// Observability + benchmark hooks.
|
||||
private val parsed = AtomicLong(0)
|
||||
private val reused = AtomicLong(0)
|
||||
@@ -117,7 +106,6 @@ class CachingEventDecoder(
|
||||
val cached = live.get(scanned.eventId) ?: previous.get(scanned.eventId)
|
||||
if (cached != null) {
|
||||
reused.incrementAndFetch()
|
||||
lastActivityAt = TimeUtils.nowMillis()
|
||||
return EventMessage(scanned.subId, cached)
|
||||
}
|
||||
}
|
||||
@@ -125,7 +113,6 @@ class CachingEventDecoder(
|
||||
val msg = fullParser.decode(text)
|
||||
if (msg is EventMessage) {
|
||||
parsed.incrementAndFetch()
|
||||
lastActivityAt = TimeUtils.nowMillis()
|
||||
if (live.size() >= capacity) {
|
||||
// Benign-race rotation: see class kdoc.
|
||||
previous = live
|
||||
@@ -137,29 +124,35 @@ class CachingEventDecoder(
|
||||
}
|
||||
|
||||
/**
|
||||
* Drops both generations once no EVENT frame has arrived for [idleMillis].
|
||||
* Retires the live generation, so entries age out on the caller's clock instead of
|
||||
* only when the cache fills.
|
||||
*
|
||||
* [capacity] bounds what the cache costs *during* a burst; this bounds how long
|
||||
* it costs anything *after* one. Without it the maps only ever rotate inside an
|
||||
* insert, so a client that goes quiet pins up to `2 * capacity` events for the
|
||||
* rest of the process's life — memory whose dedup value has already decayed to
|
||||
* zero.
|
||||
* [capacity] bounds what the cache costs *during* a burst; this bounds how long it
|
||||
* costs anything *after* one. Without it the maps only rotate inside an insert, so
|
||||
* a client receiving a slow trickle keeps everything it saw during the initial
|
||||
* burst forever — memory whose dedup value is long gone.
|
||||
*
|
||||
* Racy by the same design as rotation (see class kdoc): clearing concurrently
|
||||
* with an insert can only lose an id, and losing an id costs a re-parse, never
|
||||
* a wrong message. [nowMillis] is injectable so tests need no clock.
|
||||
* An id survives one call and dies on the next, so ticking every 30s keeps nothing
|
||||
* older than ~60s. Two consecutive calls with no traffic in between leave the cache
|
||||
* empty, which is what releases memory when the client really does go quiet.
|
||||
*
|
||||
* Racy by the same design as capacity rotation (see class kdoc): retiring a
|
||||
* generation concurrently with an insert can only lose an id, and a lost id costs a
|
||||
* re-parse, never a wrong message.
|
||||
*/
|
||||
override fun trimIfIdle(
|
||||
idleMillis: Long,
|
||||
nowMillis: Long,
|
||||
): Boolean {
|
||||
if (nowMillis - lastActivityAt < idleMillis) return false
|
||||
if (live.size() == 0 && previous.size() == 0) return false
|
||||
override fun ageOutCache() {
|
||||
previous = live
|
||||
live = ConcurrentHashCache()
|
||||
}
|
||||
|
||||
override fun clearCache() {
|
||||
previous = ConcurrentHashCache()
|
||||
live = ConcurrentHashCache()
|
||||
return true
|
||||
}
|
||||
|
||||
/** Ids currently addressable across both generations. Observability + tests. */
|
||||
val cachedCount: Int get() = live.size() + previous.size()
|
||||
|
||||
private class ScannedEvent(
|
||||
val subId: String,
|
||||
val eventId: HexKey,
|
||||
|
||||
+12
-12
@@ -21,7 +21,6 @@
|
||||
package com.vitorpamplona.quartz.nip01Core.relay.commands.toClient
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.core.OptimizedJsonMapper
|
||||
import com.vitorpamplona.quartz.utils.TimeUtils
|
||||
|
||||
/**
|
||||
* Turns one raw relay frame into a [Message]. The strategy the relay client
|
||||
@@ -35,20 +34,21 @@ fun interface MessageDecoder {
|
||||
fun decode(text: String): Message
|
||||
|
||||
/**
|
||||
* Releases whatever the decoder is holding if it has seen no frames for
|
||||
* [idleMillis]. Returns true when something was released.
|
||||
* Ages out whatever the decoder is holding. Call periodically: an entry must
|
||||
* survive at most two calls, so a 30s tick keeps nothing older than ~60s.
|
||||
*
|
||||
* A caching decoder's value decays with time — a duplicate that has not
|
||||
* arrived in a minute is not going to — but its memory does not, so the
|
||||
* owner ticks this to let a quiet decoder drop what it cached. Stateless
|
||||
* decoders keep the no-op default.
|
||||
* A caching decoder's value decays with time — a duplicate that has not arrived
|
||||
* within seconds is not going to — while its memory does not decay at all. Aging
|
||||
* on a clock rather than on idleness is deliberate: relays keep pushing a trickle
|
||||
* of events down open subscriptions forever, so a decoder is never idle in the
|
||||
* sense of "saw no frames", and an idle-triggered release would never fire.
|
||||
*
|
||||
* Never affects correctness: releasing a cache can only cost a re-parse.
|
||||
* Never affects correctness: dropping a cached id can only cost a re-parse.
|
||||
*/
|
||||
fun trimIfIdle(
|
||||
idleMillis: Long,
|
||||
nowMillis: Long = TimeUtils.nowMillis(),
|
||||
): Boolean = false
|
||||
fun ageOutCache() = Unit
|
||||
|
||||
/** Releases everything cached, e.g. when the host is shutting the client down. */
|
||||
fun clearCache() = Unit
|
||||
|
||||
companion object {
|
||||
/** Plain full-JSON parse of every frame. */
|
||||
|
||||
+67
-60
@@ -21,25 +21,24 @@
|
||||
package com.vitorpamplona.quartz.nip01Core.relay.commands.toClient
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.utils.TimeUtils
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFalse
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
/**
|
||||
* [CachingEventDecoder.trimIfIdle] bounds how long the cache costs memory, which
|
||||
* [CachingEventDecoder.capacity] does not: the generations only rotate inside an
|
||||
* insert, so without a trim a client that goes quiet pins every cached event for
|
||||
* the rest of the process's life.
|
||||
* [CachingEventDecoder.ageOutCache] bounds how long the cache costs memory, which
|
||||
* [CachingEventDecoder.capacity] does not: the generations otherwise only rotate
|
||||
* inside an insert, so a client that drops to a trickle keeps everything from its
|
||||
* initial burst forever.
|
||||
*
|
||||
* The rules that matter here are "release when genuinely idle", "do NOT release
|
||||
* out from under live traffic" — including traffic that is *only* cache hits,
|
||||
* which is exactly when the cache is earning its keep — and "releasing never
|
||||
* changes what a decode returns, only whether it re-parsed".
|
||||
* Aging is on a clock, not on idleness, and that distinction is load-bearing —
|
||||
* measured on device, relays keep pushing events down open subscriptions
|
||||
* indefinitely, so an idle-triggered release returned false on every tick across
|
||||
* 240s of a flat heap. The rules pinned here are the lifetime guarantee ("an entry
|
||||
* survives one call and dies on the next"), that traffic between calls is kept, and
|
||||
* that aging only ever costs a re-parse — never a different message.
|
||||
*
|
||||
* In `commonTest`: this is `commonMain` logic with no platform surface, and the
|
||||
* clock is injected so no target needs a real one.
|
||||
* In `commonTest`: `commonMain` logic, no platform surface, no clock.
|
||||
*/
|
||||
class CachingEventDecoderTrimTest {
|
||||
private fun hexId(seed: Int) = seed.toString(16).padStart(64, '0')
|
||||
@@ -60,78 +59,84 @@ class CachingEventDecoderTrimTest {
|
||||
subId: String = "s",
|
||||
) = """["EVENT","$subId",${event(seed).toJson()}]"""
|
||||
|
||||
private val minute = 60_000L
|
||||
|
||||
@Test
|
||||
fun releasesTheCacheOnceIdle() {
|
||||
fun entriesSurviveExactlyOneAgeOut() {
|
||||
val decoder = CachingEventDecoder(capacity = 1_000)
|
||||
(1..50).forEach { decoder.decode(frame(it)) }
|
||||
assertEquals(50, decoder.parsedCount.toInt())
|
||||
|
||||
// read the clock AFTER the decodes: they set the idle clock as they run, so a
|
||||
// cutoff measured from before them is only a full minute later while the
|
||||
// parsing itself takes under a millisecond
|
||||
val quietSince = TimeUtils.nowMillis()
|
||||
assertTrue(decoder.trimIfIdle(minute, quietSince + minute), "should trim after a quiet minute")
|
||||
|
||||
// the ids are gone, so a repeat of a previously cached frame parses again
|
||||
decoder.decode(frame(1))
|
||||
assertEquals(51, decoder.parsedCount.toInt())
|
||||
assertEquals(0, decoder.reusedCount.toInt())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun keepsTheCacheWhileTrafficIsFlowing() {
|
||||
val decoder = CachingEventDecoder(capacity = 1_000)
|
||||
val t0 = TimeUtils.nowMillis()
|
||||
(1..50).forEach { decoder.decode(frame(it)) }
|
||||
|
||||
assertFalse(decoder.trimIfIdle(minute, t0 + 1_000), "1s of quiet is not idle")
|
||||
assertEquals(1, decoder.parsedCount.toInt())
|
||||
|
||||
// first tick: retired into `previous`, still addressable
|
||||
decoder.ageOutCache()
|
||||
decoder.decode(frame(1))
|
||||
assertEquals(50, decoder.parsedCount.toInt(), "cached ids must survive")
|
||||
assertEquals(1, decoder.parsedCount.toInt(), "one age-out must not drop the id")
|
||||
assertEquals(1, decoder.reusedCount.toInt())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun cacheHitsCountAsActivity() {
|
||||
// A stream of pure duplicates inserts nothing, but it is precisely when the
|
||||
// cache is most valuable. Keying idleness off inserts alone would trim it away
|
||||
// under the traffic it is serving.
|
||||
fun twoAgeOutsWithNoTrafficEmptyTheCache() {
|
||||
val decoder = CachingEventDecoder(capacity = 1_000)
|
||||
(1..50).forEach { decoder.decode(frame(it)) }
|
||||
assertEquals(50, decoder.cachedCount)
|
||||
|
||||
decoder.ageOutCache()
|
||||
decoder.ageOutCache()
|
||||
|
||||
assertEquals(0, decoder.cachedCount, "nothing survives two ticks without traffic")
|
||||
decoder.decode(frame(1))
|
||||
assertEquals(51, decoder.parsedCount.toInt(), "the id is gone, so it re-parses")
|
||||
assertEquals(0, decoder.reusedCount.toInt())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun trafficBetweenTicksIsKept() {
|
||||
// the trickle case: events keep arriving, so the cache must keep serving them
|
||||
val decoder = CachingEventDecoder(capacity = 1_000)
|
||||
decoder.decode(frame(1))
|
||||
decoder.ageOutCache()
|
||||
decoder.decode(frame(2))
|
||||
decoder.ageOutCache()
|
||||
|
||||
val muchLater = TimeUtils.nowMillis() + 10 * minute
|
||||
repeat(5) { decoder.decode(frame(1, subId = "sub$it")) }
|
||||
// 1 was inserted two ticks ago -> gone; 2 was inserted one tick ago -> kept
|
||||
decoder.decode(frame(2))
|
||||
assertEquals(2, decoder.parsedCount.toInt())
|
||||
assertEquals(1, decoder.reusedCount.toInt(), "the recent event is still cached")
|
||||
|
||||
assertEquals(5, decoder.reusedCount.toInt())
|
||||
assertFalse(
|
||||
decoder.trimIfIdle(minute, TimeUtils.nowMillis() + 1_000),
|
||||
"hits must refresh the idle clock",
|
||||
)
|
||||
// ...but a genuinely quiet stretch after those hits still trims
|
||||
assertTrue(decoder.trimIfIdle(minute, muchLater))
|
||||
decoder.decode(frame(1))
|
||||
assertEquals(3, decoder.parsedCount.toInt(), "the older event aged out")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun trimmingAnEmptyCacheDoesNothing() {
|
||||
val decoder = CachingEventDecoder(capacity = 1_000)
|
||||
assertFalse(decoder.trimIfIdle(0, TimeUtils.nowMillis() + minute), "nothing to release")
|
||||
fun agingABusyCacheKeepsItSmallInsteadOfUnbounded() {
|
||||
// what the device measurement is about: a burst, then a trickle. Without aging
|
||||
// the burst's ids stay resident forever.
|
||||
val decoder = CachingEventDecoder(capacity = 100_000)
|
||||
(1..5_000).forEach { decoder.decode(frame(it)) }
|
||||
assertEquals(5_000, decoder.cachedCount)
|
||||
|
||||
decoder.ageOutCache()
|
||||
(5_001..5_010).forEach { decoder.decode(frame(it)) } // the trickle
|
||||
decoder.ageOutCache()
|
||||
|
||||
assertEquals(10, decoder.cachedCount, "only the trickle survives, not the burst")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun zeroIdleReleasesImmediately() {
|
||||
fun clearCacheReleasesEverythingAtOnce() {
|
||||
// what NostrClient.disconnect() uses when the host backgrounds the app
|
||||
val decoder = CachingEventDecoder(capacity = 1_000)
|
||||
(1..10).forEach { decoder.decode(frame(it)) }
|
||||
assertTrue(decoder.trimIfIdle(idleMillis = 0, nowMillis = TimeUtils.nowMillis()))
|
||||
assertTrue(decoder.cachedCount > 0)
|
||||
|
||||
decoder.clearCache()
|
||||
assertEquals(0, decoder.cachedCount)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun aTrimNeverChangesWhatADecodeReturns() {
|
||||
fun agingNeverChangesWhatADecodeReturns() {
|
||||
val decoder = CachingEventDecoder(capacity = 1_000)
|
||||
val before = decoder.decode(frame(7, subId = "a")) as EventMessage
|
||||
decoder.trimIfIdle(0, TimeUtils.nowMillis())
|
||||
decoder.ageOutCache()
|
||||
decoder.ageOutCache()
|
||||
val after = decoder.decode(frame(7, subId = "a")) as EventMessage
|
||||
|
||||
assertEquals(before.subId, after.subId)
|
||||
@@ -139,12 +144,14 @@ class CachingEventDecoderTrimTest {
|
||||
assertEquals(before.event.pubKey, after.event.pubKey)
|
||||
assertEquals(before.event.content, after.event.content)
|
||||
assertEquals(before.event.createdAt, after.event.createdAt)
|
||||
// only the parse count differs: the trim cost one re-parse, nothing else
|
||||
// only the parse count differs: aging cost one re-parse, nothing else
|
||||
assertEquals(2, decoder.parsedCount.toInt())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun theDefaultDecoderIgnoresTrim() {
|
||||
assertFalse(MessageDecoder.Default.trimIfIdle(0, TimeUtils.nowMillis()), "stateless: nothing to release")
|
||||
fun theDefaultDecoderIgnoresBoth() {
|
||||
// stateless: nothing to age or release, and neither call may throw
|
||||
MessageDecoder.Default.ageOutCache()
|
||||
MessageDecoder.Default.clearCache()
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user