feat: paginate relay queries with per-relay until cursor
Relays cap event responses at ~500 per request. After each EOSE, the oldest createdAt seen on a relay becomes the next 'until' cursor and the relay is re-queried for the next page. This repeats until a relay returns no events (exhausted) or cannot be reached. All active relays in a batch are queried in parallel each round; only the relays that still have pages remaining participate in the next round. Dead relays (onCannotConnect) are dropped immediately. ConcurrentHashMap is used for per-relay counters and cursors since IRequestListener callbacks can arrive from concurrent relay threads. https://claude.ai/code/session_01U8qF9mK4UBvXsXP1zNMXfX
This commit is contained in:
+86
-43
@@ -36,6 +36,7 @@ import kotlinx.coroutines.flow.StateFlow
|
|||||||
import kotlinx.coroutines.isActive
|
import kotlinx.coroutines.isActive
|
||||||
import kotlinx.coroutines.launch
|
import kotlinx.coroutines.launch
|
||||||
import kotlinx.coroutines.withTimeoutOrNull
|
import kotlinx.coroutines.withTimeoutOrNull
|
||||||
|
import java.util.concurrent.ConcurrentHashMap
|
||||||
import java.util.concurrent.atomic.AtomicInteger
|
import java.util.concurrent.atomic.AtomicInteger
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -245,62 +246,104 @@ class EventSyncViewModel(
|
|||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Sends [filters] to every relay in [relays] simultaneously (one subscription covering
|
* Sends [baseFilters] to every relay in [relays] simultaneously, then paginates per relay
|
||||||
* the whole chunk). Suspends until every relay has replied with EOSE/CLOSED/error, or
|
* until each one is exhausted.
|
||||||
* [CHUNK_TIMEOUT_MS] elapses, whichever comes first.
|
*
|
||||||
|
* Many relays impose a hard event-count limit per filter (commonly 500). After EOSE the
|
||||||
|
* oldest [Event.createdAt] seen on that relay becomes the next `until` cursor, so the
|
||||||
|
* following request fetches the next page. This repeats until a relay returns no new
|
||||||
|
* events, at which point it is removed from the active set.
|
||||||
|
*
|
||||||
|
* All relays in [relays] are queried in parallel within each pagination round. Only the
|
||||||
|
* relays that still have more pages participate in subsequent rounds.
|
||||||
*/
|
*/
|
||||||
private suspend fun downloadFromChunk(
|
private suspend fun downloadFromChunk(
|
||||||
relays: List<NormalizedRelayUrl>,
|
relays: List<NormalizedRelayUrl>,
|
||||||
filters: List<Filter>,
|
baseFilters: List<Filter>,
|
||||||
onEvent: (Event) -> Unit,
|
onEvent: (Event) -> Unit,
|
||||||
) {
|
) {
|
||||||
val remaining = AtomicInteger(relays.size)
|
// `until` cursor per relay; absent = first page (no cursor).
|
||||||
val done = Channel<Unit>(Channel.CONFLATED)
|
val relayUntil = ConcurrentHashMap<NormalizedRelayUrl, Long>()
|
||||||
val subId = newSubId()
|
// Relays that can no longer be contacted are skipped immediately.
|
||||||
|
val deadRelays = ConcurrentHashMap.newKeySet<NormalizedRelayUrl>()
|
||||||
|
|
||||||
val listener =
|
val activeRelays = relays.toMutableList()
|
||||||
object : IRequestListener {
|
|
||||||
override fun onEvent(
|
while (activeRelays.isNotEmpty() && isActive) {
|
||||||
event: Event,
|
// Per-relay event count and oldest timestamp for this round.
|
||||||
isLive: Boolean,
|
val roundCount = ConcurrentHashMap<NormalizedRelayUrl, Int>()
|
||||||
relay: NormalizedRelayUrl,
|
val roundMinTs = ConcurrentHashMap<NormalizedRelayUrl, Long>()
|
||||||
forFilters: List<Filter>?,
|
|
||||||
) {
|
val remaining = AtomicInteger(activeRelays.size)
|
||||||
onEvent(event)
|
val done = Channel<Unit>(Channel.CONFLATED)
|
||||||
|
val subId = newSubId()
|
||||||
|
|
||||||
|
val listener =
|
||||||
|
object : IRequestListener {
|
||||||
|
override fun onEvent(
|
||||||
|
event: Event,
|
||||||
|
isLive: Boolean,
|
||||||
|
relay: NormalizedRelayUrl,
|
||||||
|
forFilters: List<Filter>?,
|
||||||
|
) {
|
||||||
|
onEvent(event)
|
||||||
|
roundCount.merge(relay, 1, Int::plus)
|
||||||
|
roundMinTs.merge(relay, event.createdAt, ::minOf)
|
||||||
|
}
|
||||||
|
|
||||||
|
override fun onEose(
|
||||||
|
relay: NormalizedRelayUrl,
|
||||||
|
forFilters: List<Filter>?,
|
||||||
|
) {
|
||||||
|
if (remaining.decrementAndGet() <= 0) done.trySend(Unit)
|
||||||
|
}
|
||||||
|
|
||||||
|
override fun onClosed(
|
||||||
|
message: String,
|
||||||
|
relay: NormalizedRelayUrl,
|
||||||
|
forFilters: List<Filter>?,
|
||||||
|
) {
|
||||||
|
if (remaining.decrementAndGet() <= 0) done.trySend(Unit)
|
||||||
|
}
|
||||||
|
|
||||||
|
override fun onCannotConnect(
|
||||||
|
relay: NormalizedRelayUrl,
|
||||||
|
message: String,
|
||||||
|
forFilters: List<Filter>?,
|
||||||
|
) {
|
||||||
|
deadRelays.add(relay)
|
||||||
|
if (remaining.decrementAndGet() <= 0) done.trySend(Unit)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override fun onEose(
|
// Build per-relay filter list, injecting each relay's current `until` cursor.
|
||||||
relay: NormalizedRelayUrl,
|
val filtersMap =
|
||||||
forFilters: List<Filter>?,
|
activeRelays.associateWith { relay ->
|
||||||
) {
|
val until = relayUntil[relay]
|
||||||
if (remaining.decrementAndGet() <= 0) done.trySend(Unit)
|
if (until == null) {
|
||||||
|
baseFilters
|
||||||
|
} else {
|
||||||
|
baseFilters.map { it.copy(until = until) }
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override fun onClosed(
|
account.client.openReqSubscription(subId, filtersMap, listener)
|
||||||
message: String,
|
withTimeoutOrNull(CHUNK_TIMEOUT_MS) { done.receive() }
|
||||||
relay: NormalizedRelayUrl,
|
account.client.close(subId)
|
||||||
forFilters: List<Filter>?,
|
done.close()
|
||||||
) {
|
|
||||||
if (remaining.decrementAndGet() <= 0) done.trySend(Unit)
|
|
||||||
}
|
|
||||||
|
|
||||||
override fun onCannotConnect(
|
// Advance cursors for relays that returned events; drop exhausted or dead ones.
|
||||||
relay: NormalizedRelayUrl,
|
activeRelays.removeAll { relay ->
|
||||||
message: String,
|
when {
|
||||||
forFilters: List<Filter>?,
|
relay in deadRelays -> true
|
||||||
) {
|
(roundCount[relay] ?: 0) == 0 -> true // EOSE with no events = exhausted
|
||||||
if (remaining.decrementAndGet() <= 0) done.trySend(Unit)
|
else -> {
|
||||||
|
// Next page starts just before the oldest event seen this round.
|
||||||
|
relayUntil[relay] = roundMinTs[relay]!! - 1
|
||||||
|
false
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
val filtersMap = relays.associateWith { filters }
|
|
||||||
account.client.openReqSubscription(subId, filtersMap, listener)
|
|
||||||
|
|
||||||
withTimeoutOrNull(CHUNK_TIMEOUT_MS) {
|
|
||||||
done.receive()
|
|
||||||
}
|
}
|
||||||
|
|
||||||
account.client.close(subId)
|
|
||||||
done.close()
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user