mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 19:53:08 +00:00
revert(relay): drop FrameDispatchStats — per-frame cost on the shared WS hot path
FrameDispatchStats stamped a ValueTimeMark on every relay frame and recorded a contended atomic per frame in BasicOkHttpWebSocket — the WebSocket layer used by the whole app, unconditionally, forever — to answer a one-time question that only graperank --diagnose read. It served its purpose (proved the our-side dispatch lag is ~200ms mean and the EOSE-wait is dominantly relay-side, so the crawler is network-bound), but the ongoing per-frame Pair allocation + atomic contention on every client's relay traffic isn't worth carrying. Revert the channel back to Channel<String> and delete the stats holder. The diagnose-gated saturation ticker and per-drain latency breakdown stay — they're crawler-local, off the hot path. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01MSW59hJtP4Yn8fnRUxc7F5
This commit is contained in:
-6
@@ -26,7 +26,6 @@ import com.vitorpamplona.quartz.nip01Core.crypto.verify
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.AdaptiveRelayLimiter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.DrainFailure
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.FrameDispatchStats
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.classifyDrainFailure
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAllPages
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.reqs.SubscriptionListener
|
||||
@@ -257,7 +256,6 @@ class GrapeRankCrawler(
|
||||
verifyNanos.store(0)
|
||||
insertNanos.store(0)
|
||||
eventsStored.store(0)
|
||||
if (config.diagnose) FrameDispatchStats.reset()
|
||||
return CrawlRun(observer, builder).run()
|
||||
}
|
||||
|
||||
@@ -1597,10 +1595,6 @@ class GrapeRankCrawler(
|
||||
"${throttled.load()} rate-limit responses. " +
|
||||
"High EOSE-wait % → a shorter/adaptive fast window is the lever, not more concurrency.",
|
||||
)
|
||||
// Attribution: is the EOSE-wait the relay (slow to SEND eose) or us (the
|
||||
// eose frame arrived but sat in our IO pipeline)? Low dispatch-lag + high
|
||||
// EOSE-wait ⇒ relay; high dispatch-lag ⇒ our Dispatchers.IO is backed up.
|
||||
log("[graperank] frame dispatch (our-side pipeline lag): ${FrameDispatchStats.snapshot()}")
|
||||
}
|
||||
return Stats(
|
||||
rounds = rounds,
|
||||
|
||||
-84
@@ -1,84 +0,0 @@
|
||||
/*
|
||||
* 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.quartz.nip01Core.relay.client.accessories
|
||||
|
||||
import kotlin.concurrent.atomics.AtomicLong
|
||||
import kotlin.concurrent.atomics.ExperimentalAtomicApi
|
||||
|
||||
/**
|
||||
* Diagnostic: the lag between a raw relay frame arriving on the socket (the OkHttp
|
||||
* reader thread's `onMessage`) and our per-connection consumer coroutine actually
|
||||
* pulling it off the channel to decode + dispatch. Because that consumer runs on the
|
||||
* shared `Dispatchers.IO`, this lag is precisely OUR-side pipeline delay — channel
|
||||
* queue wait + coroutine reschedule + time spent decoding earlier frames — with the
|
||||
* relay's own send timing excluded (the reader thread enqueues the instant bytes land).
|
||||
*
|
||||
* It exists to answer one question the [GrapeRankCrawler]'s EOSE-wait metric cannot on
|
||||
* its own: when a drain shows a 5-second gap between the relay's last event and its
|
||||
* EOSE, is the relay slow to send EOSE (low dispatch lag) or is our IO pipeline backed
|
||||
* up so the already-arrived EOSE frame sits queued (high dispatch lag)? Process-global
|
||||
* and opt-in: a caller [reset]s before a run and reads [snapshot] after. Off unless
|
||||
* something records into it, so zero cost on normal paths.
|
||||
*/
|
||||
@OptIn(ExperimentalAtomicApi::class)
|
||||
object FrameDispatchStats {
|
||||
private val count = AtomicLong(0)
|
||||
private val sumMs = AtomicLong(0)
|
||||
private val maxMs = AtomicLong(0)
|
||||
private val over1s = AtomicLong(0)
|
||||
|
||||
fun record(lagMs: Long) {
|
||||
count.addAndFetch(1)
|
||||
sumMs.addAndFetch(lagMs)
|
||||
if (lagMs >= 1000) over1s.addAndFetch(1)
|
||||
// Lock-free running max.
|
||||
while (true) {
|
||||
val cur = maxMs.load()
|
||||
if (lagMs <= cur || maxMs.compareAndSet(cur, lagMs)) break
|
||||
}
|
||||
}
|
||||
|
||||
fun reset() {
|
||||
count.store(0)
|
||||
sumMs.store(0)
|
||||
maxMs.store(0)
|
||||
over1s.store(0)
|
||||
}
|
||||
|
||||
class Snapshot(
|
||||
val frames: Long,
|
||||
val meanMs: Long,
|
||||
val maxMs: Long,
|
||||
val over1s: Long,
|
||||
) {
|
||||
override fun toString() = "frames=$frames, mean dispatch-lag=${meanMs}ms, max=${maxMs}ms, $over1s frames waited >1s in our pipeline"
|
||||
}
|
||||
|
||||
fun snapshot(): Snapshot {
|
||||
val n = count.load()
|
||||
return Snapshot(
|
||||
frames = n,
|
||||
meanMs = if (n > 0) sumMs.load() / n else 0,
|
||||
maxMs = maxMs.load(),
|
||||
over1s = over1s.load(),
|
||||
)
|
||||
}
|
||||
}
|
||||
+4
-13
@@ -20,7 +20,6 @@
|
||||
*/
|
||||
package com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.FrameDispatchStats
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebSocket
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebSocketListener
|
||||
@@ -36,8 +35,6 @@ import kotlinx.coroutines.launch
|
||||
import okhttp3.OkHttpClient
|
||||
import okhttp3.Request
|
||||
import okhttp3.Response
|
||||
import kotlin.time.TimeSource
|
||||
import kotlin.time.TimeSource.Monotonic.ValueTimeMark
|
||||
import okhttp3.WebSocket as OkHttpWebSocket
|
||||
import okhttp3.WebSocketListener as OkHttpWebSocketListener
|
||||
|
||||
@@ -74,14 +71,10 @@ class BasicOkHttpWebSocket(
|
||||
// fast as it can send and own the buffering; consumer speed
|
||||
// is handled downstream (CachingEventDecoder,
|
||||
// ParallelEventVerifier).
|
||||
val incomingMessages: Channel<Pair<ValueTimeMark, String>> = Channel(Channel.UNLIMITED)
|
||||
val incomingMessages: Channel<String> = Channel(Channel.UNLIMITED)
|
||||
val job = // Launch a coroutine to process messages from the channel.
|
||||
scope.launch {
|
||||
for ((arrivedAt, message) in incomingMessages) {
|
||||
// Lag from raw socket arrival to us pulling it off the channel =
|
||||
// OUR-side pipeline delay (queue wait + IO reschedule + decoding
|
||||
// earlier frames), with the relay's send timing excluded.
|
||||
FrameDispatchStats.record(arrivedAt.elapsedNow().inWholeMilliseconds)
|
||||
for (message in incomingMessages) {
|
||||
out.onMessage(message)
|
||||
}
|
||||
}
|
||||
@@ -99,10 +92,8 @@ class BasicOkHttpWebSocket(
|
||||
text: String,
|
||||
) {
|
||||
// Never blocks (unlimited channel): the OkHttp reader
|
||||
// thread must stay free to keep draining the socket. Stamp the
|
||||
// arrival instant here (on the reader thread, before any of our
|
||||
// queueing) so downstream dispatch lag is measurable.
|
||||
incomingMessages.trySendBlocking(TimeSource.Monotonic.markNow() to text)
|
||||
// thread must stay free to keep draining the socket.
|
||||
incomingMessages.trySendBlocking(text)
|
||||
}
|
||||
|
||||
override fun onClosed(
|
||||
|
||||
Reference in New Issue
Block a user