mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-08-12 09:13:23 +00:00
perf(graperank): idle-based park timeout so streaming relays aren't cut mid-flight
The park window's timeout was absolute from subscription open, so a relay still actively streaming a large result set once it passed parkTimeoutMs was unsubscribed and its untransmitted tail lost. Reset the window on every incoming event (a conflated activity signal drives a select against the terminal deferred), so a parked subscription is closed only after parkTimeoutMs of actual silence — never while events are still arriving. The fast window stays absolute: it only decides when to hand a slow relay to the background park lane, which loses nothing. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01RWk2ZMrGBSr4WenKgwqmbB
This commit is contained in:
+43
-1
@@ -50,6 +50,7 @@ import kotlinx.coroutines.channels.Channel
|
||||
import kotlinx.coroutines.coroutineScope
|
||||
import kotlinx.coroutines.joinAll
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.selects.select
|
||||
import kotlinx.coroutines.withTimeoutOrNull
|
||||
import kotlin.concurrent.atomics.AtomicLong
|
||||
import kotlin.concurrent.atomics.ExperimentalAtomicApi
|
||||
@@ -559,6 +560,34 @@ class GrapeRankDataCrawler(
|
||||
return fresh
|
||||
}
|
||||
|
||||
/**
|
||||
* Wait for a subscription's terminal ([done]: EOSE/CLOSED/cannot), resetting
|
||||
* the [idleMs] window every time an event pings [activity]. So the wait ends
|
||||
* with "timeout" only after [idleMs] of actual SILENCE — a relay that keeps
|
||||
* streaming (however long its result set) is never cut mid-flight; only a
|
||||
* genuinely stalled one is. Used for the patient park window.
|
||||
*/
|
||||
private suspend fun awaitTerminalOrIdle(
|
||||
done: CompletableDeferred<String>,
|
||||
activity: Channel<Unit>,
|
||||
idleMs: Long,
|
||||
): String {
|
||||
while (true) {
|
||||
val r =
|
||||
withTimeoutOrNull(idleMs) {
|
||||
select {
|
||||
done.onAwait { it }
|
||||
activity.onReceive { ACTIVITY }
|
||||
}
|
||||
}
|
||||
when (r) {
|
||||
null -> return "timeout" // idleMs elapsed with no event and no terminal
|
||||
ACTIVITY -> Unit // an event arrived — reset the idle window and keep waiting
|
||||
else -> return r // terminal reason
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Fold one late-delivered event from a parked relay into the graph. Only the
|
||||
* round loop calls this (directly or via [foldLateHarvest]), so graph state
|
||||
@@ -677,6 +706,10 @@ class GrapeRankDataCrawler(
|
||||
val subId = newSubId()
|
||||
val done = CompletableDeferred<String>()
|
||||
val unitEvents = Channel<Pair<NormalizedRelayUrl, Event>>(Channel.UNLIMITED)
|
||||
// Liveness signal for the parked idle timeout: every event pings
|
||||
// this (conflated, so bursts collapse to one) and resets the park
|
||||
// window, so a relay actively streaming is never cut mid-flight.
|
||||
val activity = Channel<Unit>(Channel.CONFLATED)
|
||||
val listener =
|
||||
object : SubscriptionListener {
|
||||
override fun onEvent(
|
||||
@@ -686,6 +719,7 @@ class GrapeRankDataCrawler(
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
unitEvents.trySend(relay to event)
|
||||
activity.trySend(Unit)
|
||||
}
|
||||
|
||||
override fun onEose(
|
||||
@@ -732,7 +766,10 @@ class GrapeRankDataCrawler(
|
||||
parkedInFlight.addAndFetch(1)
|
||||
scope.launch {
|
||||
try {
|
||||
val late = withTimeoutOrNull(config.parkTimeoutMs) { done.await() } ?: "timeout"
|
||||
// Idle timeout, not absolute: only cut after parkTimeoutMs
|
||||
// of SILENCE (no event, no terminal), so a relay still
|
||||
// streaming a large result set is never chopped mid-flight.
|
||||
val late = awaitTerminalOrIdle(done, activity, config.parkTimeoutMs)
|
||||
logSlow(subRelay, "parked→$late", mark.elapsedNow().inWholeMilliseconds, groupFilters)
|
||||
// A parked relay that ends in a hard/transient failure (not a
|
||||
// clean EOSE) is reported dead the same way a fast one would be.
|
||||
@@ -997,6 +1034,11 @@ class GrapeRankDataCrawler(
|
||||
// to block waiting for one of them to deliver before re-checking convergence.
|
||||
private const val PARK_POLL_MS = 2000L
|
||||
|
||||
// Sentinel returned by the park idle-wait's select when an event arrived
|
||||
// (resets the window). A control string that can't collide with a relay's
|
||||
// CLOSED/cannot message, which are the only other select results.
|
||||
private const val ACTIVITY = " | ||||