mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-08-10 00:16:59 +00:00
refactor(dm): page rooms-list NIP-04 history per relay; drop dead pager code
Convert ChatroomListNip04HistorySubAssembler from the round-based model to the per-relay-independent one, matching the per-conversation NIP-04 loader and the gift-wrap history loader. A single loadMore opens one window spanning the whole walk; each relay continues itself on its own non-empty EOSE (beginRound([relay]) + invalidateFilters), and silent/unreachable relays are marked stalled (kept open) with exhaustion coming from the WindowLoadTracker's silence + connect-grace backstops (tracksReqSends=true). The rooms list now pages both protocols identically. This was the last user of the round-tally + give-up machinery, so remove it from UntilLimitPager: roundEventCount(), onClosed()/giveUp() and the givenUp/closedStreak cursor state + GIVE_UP_AFTER_CLOSES. activeRelays now excludes only done relays. Delete UntilLimitPagerGiveUpTest (tested the removed give-up path; abandonment is now the WindowLoadTracker's job, which its own test covers). Also drop the now-dead loadEverything/autoLoadAll and the no-progress retry (no callers; with independent paging one loadMore already walks each relay to its bottom, and the tracker's backstops cover cold starts).
This commit is contained in:
+6
-69
@@ -35,9 +35,10 @@ import java.util.concurrent.ConcurrentHashMap
|
||||
* Stop signal (per relay): an **empty page followed by EOSE** ([onEose] with no events) marks that
|
||||
* relay [done][isDone]. A relay that returns anything — even fewer than the requested limit, since a
|
||||
* relay may cap results below what we asked — is *not* done; its cursor advances to one second below
|
||||
* the oldest event it sent and it is asked again. Globally, the owner decides "exhausted" from
|
||||
* [roundEventCount]: a whole round that advanced no relay (every relay empty-EOSE'd or only answered
|
||||
* CLOSED) means nothing more is reachable.
|
||||
* the oldest event it sent and the owner asks it again (typically right away, per relay). Relays a
|
||||
* relay can't be read from (CLOSED / unreachable / silent) are the owner's concern, not the cursor's:
|
||||
* the owner flags them stalled and lets its [WindowLoadTracker] settle the load so they can't block
|
||||
* exhaustion forever.
|
||||
*
|
||||
* Not internally synchronized: per-relay counters are touched on the relay IO threads (one relay's
|
||||
* callbacks are serialized) and read on the owning scope after the load settles; fields are volatile.
|
||||
@@ -50,14 +51,6 @@ class UntilLimitPager<K> {
|
||||
// Set once the relay answered an empty page with EOSE: there is nothing older on it.
|
||||
@Volatile var done: Boolean = false
|
||||
|
||||
// Set once the relay has rejected us [GIVE_UP_AFTER_CLOSES] times in a row without ever
|
||||
// answering (CLOSED, e.g. "auth-required" for authors we can't authenticate). It is NOT done —
|
||||
// we just can't read its window — but it must be excluded so it doesn't block exhaustion forever.
|
||||
@Volatile var givenUp: Boolean = false
|
||||
|
||||
// Consecutive CLOSEDs since this relay last actually answered (event or EOSE). Reset on contact.
|
||||
@Volatile var closedStreak: Int = 0
|
||||
|
||||
// Per-round tallies, reset by [beginRound]: how many events arrived and the oldest among them.
|
||||
@Volatile var roundCount: Int = 0
|
||||
|
||||
@@ -103,7 +96,6 @@ class UntilLimitPager<K> {
|
||||
createdAt: Long,
|
||||
) {
|
||||
val c = cursor(key, relay)
|
||||
c.closedStreak = 0
|
||||
c.roundCount++
|
||||
if (createdAt < c.roundOldest) c.roundOldest = createdAt
|
||||
}
|
||||
@@ -118,7 +110,6 @@ class UntilLimitPager<K> {
|
||||
relay: NormalizedRelayUrl,
|
||||
) {
|
||||
val c = cursor(key, relay)
|
||||
c.closedStreak = 0
|
||||
if (c.roundCount == 0) {
|
||||
c.done = true
|
||||
} else {
|
||||
@@ -126,59 +117,11 @@ class UntilLimitPager<K> {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Records a CLOSED (rejection) from [relay]. After [GIVE_UP_AFTER_CLOSES] in a row with no answer in
|
||||
* between — i.e. the relay keeps rejecting us and auth can't fix it — the relay is [given up][givenUp]
|
||||
* so it stops blocking exhaustion. Returns true if this CLOSED tipped it into given-up.
|
||||
*/
|
||||
fun onClosed(
|
||||
key: K,
|
||||
relay: NormalizedRelayUrl,
|
||||
): Boolean {
|
||||
val c = cursor(key, relay)
|
||||
if (c.givenUp || c.done) return false
|
||||
c.closedStreak++
|
||||
if (c.closedStreak >= GIVE_UP_AFTER_CLOSES) {
|
||||
c.givenUp = true
|
||||
return true
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
/**
|
||||
* Abandons [relay] for [key] when it accepted our REQ but never answered (no event, EOSE, or CLOSED
|
||||
* within the silence window). Like [givenUp] via [onClosed], it is excluded from [activeRelays] so a
|
||||
* silent relay can't block exhaustion forever — but a relay that already finished cleanly ([done])
|
||||
* is left alone. Returns true if this abandoned a relay that wasn't already done/given-up.
|
||||
*/
|
||||
fun giveUp(
|
||||
key: K,
|
||||
relay: NormalizedRelayUrl,
|
||||
): Boolean {
|
||||
val c = cursor(key, relay)
|
||||
if (c.done || c.givenUp) return false
|
||||
c.givenUp = true
|
||||
return true
|
||||
}
|
||||
|
||||
/** Total events received across [relays] in the round just finished. Zero ⇒ nothing more is reachable. */
|
||||
fun roundEventCount(
|
||||
key: K,
|
||||
relays: Collection<NormalizedRelayUrl>,
|
||||
): Int = relays.sumOf { cursor(key, it).roundCount }
|
||||
|
||||
/**
|
||||
* Relays from [all] that still have older history to ask for: not yet empty-EOSE'd ([done]) and not
|
||||
* abandoned as unreadable ([givenUp]).
|
||||
*/
|
||||
/** Relays from [all] that still have older history to ask for: not yet empty-EOSE'd ([done]). */
|
||||
fun activeRelays(
|
||||
key: K,
|
||||
all: Collection<NormalizedRelayUrl>,
|
||||
): List<NormalizedRelayUrl> =
|
||||
all.filterNot {
|
||||
val c = cursor(key, it)
|
||||
c.done || c.givenUp
|
||||
}
|
||||
): List<NormalizedRelayUrl> = all.filterNot { cursor(key, it).done }
|
||||
|
||||
/**
|
||||
* The oldest point reached across [relays] — the minimum cursor (how far back paging has gone).
|
||||
@@ -189,10 +132,4 @@ class UntilLimitPager<K> {
|
||||
relays: Collection<NormalizedRelayUrl>,
|
||||
start: Long,
|
||||
): Long? = relays.takeIf { it.isNotEmpty() }?.minOf { cursor(key, it).until ?: start }
|
||||
|
||||
companion object {
|
||||
// Consecutive CLOSEDs (with no answer in between) before a relay is abandoned as unreadable.
|
||||
// Allows for the pool's auth handshake + a retry or two before concluding auth can't succeed.
|
||||
private const val GIVE_UP_AFTER_CLOSES = 3
|
||||
}
|
||||
}
|
||||
|
||||
+152
-116
@@ -40,7 +40,6 @@ import com.vitorpamplona.quartz.utils.Log
|
||||
import com.vitorpamplona.quartz.utils.TimeUtils
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Job
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.flow.MutableStateFlow
|
||||
import kotlinx.coroutines.flow.StateFlow
|
||||
import kotlinx.coroutines.flow.asStateFlow
|
||||
@@ -48,11 +47,20 @@ import kotlinx.coroutines.launch
|
||||
import java.util.concurrent.ConcurrentHashMap
|
||||
|
||||
/**
|
||||
* Loads older NIP-04 DMs (kind 4) for the rooms list by `until`+`limit` paging, per relay — the same
|
||||
* gap-proof model as [com.vitorpamplona.amethyst.service.relayClient.reqCommand.account.nip59GiftWraps.AccountGiftWrapsHistoryEoseManager],
|
||||
* but for kind 4 (exact timestamps, no margin) across the account's home + DM relays. Idle until
|
||||
* [loadMore]; a relay is done on an empty page + EOSE; the whole history is [exhausted] once a round
|
||||
* advances no relay.
|
||||
* Loads older NIP-04 DMs (kind 4) for the rooms list by `until`+`limit` paging, **per relay,
|
||||
* independently** — the same model as the per-conversation loader
|
||||
* ([com.vitorpamplona.amethyst.ui.screen.loggedIn.chats.privateDM.datasource.ChatroomNip04HistorySubAssembler])
|
||||
* and the gift-wrap history loader
|
||||
* ([com.vitorpamplona.amethyst.service.relayClient.reqCommand.account.nip59GiftWraps.AccountGiftWrapsHistoryEoseManager]),
|
||||
* but account-wide across the home (outbox, *from me*) + DM (inbox, *to me*) relays.
|
||||
*
|
||||
* Idle until [loadMore]. A single [loadMore] kicks off every relay that still has older history, and
|
||||
* from then on each relay drives its own pages off its own cursor, continuing the instant it EOSEs a
|
||||
* non-empty page ([onEose] → `pager.beginRound([relay])` + `invalidateFilters`; the subscription layer
|
||||
* diffs per relay, so re-issuing only re-REQs the relay whose cursor moved). A relay is *done* on an
|
||||
* empty page + EOSE; one that won't answer (auth CLOSE, unreachable, silent) is marked *stalled* but
|
||||
* kept open. The whole history is [exhausted] once the window settles — every relay done or stalled —
|
||||
* via the [WindowLoadTracker]'s silence + connect-grace backstops.
|
||||
*/
|
||||
class ChatroomListNip04HistorySubAssembler(
|
||||
client: INostrClient,
|
||||
@@ -60,25 +68,26 @@ class ChatroomListNip04HistorySubAssembler(
|
||||
) : PerUserEoseManager<ChatroomListState>(client, allKeys) {
|
||||
private val pager = UntilLimitPager<HexKey>()
|
||||
private val started = ConcurrentHashMap.newKeySet<HexKey>()
|
||||
private val askedRelays = ConcurrentHashMap<HexKey, Set<NormalizedRelayUrl>>()
|
||||
private val accounts = ConcurrentHashMap<HexKey, Account>()
|
||||
|
||||
// Shared across accounts (singleton coordinator): repoint the display flows to the active account
|
||||
// on switch instead of leaking the previous one's exhausted/mark state. Cursors live in [pager].
|
||||
// Relays currently not advancing for a user (auth CLOSE / unreachable / silent). Tracked only for the
|
||||
// logs; these relays are NOT given up — they keep their subscription and keep trying to catch up.
|
||||
private val stalledRelays = ConcurrentHashMap<HexKey, MutableSet<NormalizedRelayUrl>>()
|
||||
|
||||
// Shared across accounts (singleton coordinator): repoint the display flows to the active account on
|
||||
// switch instead of leaking the previous one's state. Cursors live in [pager].
|
||||
@Volatile
|
||||
private var activeUser: HexKey? = null
|
||||
private val exhaustedByUser = ConcurrentHashMap<HexKey, Boolean>()
|
||||
|
||||
// No-progress guard: skip re-issuing an identical round that brought nothing (see the gift-wrap
|
||||
// history manager's twin); cleared by onEose.
|
||||
@Volatile
|
||||
private var lastAskedActive: Set<NormalizedRelayUrl> = emptySet()
|
||||
private val windowLoad = WindowLoadTracker("rooms.nip04.history", tracksReqSends = true, onAbandoned = ::onRelaysStalled)
|
||||
|
||||
@Volatile
|
||||
private var lastRoundEventCount = -1
|
||||
|
||||
private val windowLoad = WindowLoadTracker("rooms.nip04.history")
|
||||
val loadingMore: StateFlow<Boolean> = windowLoad.loading
|
||||
// Exposed instead of windowLoad.loading directly: that flow starts `true` (it assumes a load is in
|
||||
// flight from construction). Wired straight through, its `true` would wedge the scroll-driven loader
|
||||
// — whose gate is `!loading` — so the first loadMore could never fire. This starts false and only
|
||||
// goes true once paging actually begins (mirrored from windowLoad by the done collector).
|
||||
private val _loadingMore = MutableStateFlow(false)
|
||||
val loadingMore: StateFlow<Boolean> = _loadingMore.asStateFlow()
|
||||
|
||||
private val _exhausted = MutableStateFlow(false)
|
||||
val exhausted: StateFlow<Boolean> = _exhausted.asStateFlow()
|
||||
@@ -93,21 +102,30 @@ class ChatroomListNip04HistorySubAssembler(
|
||||
private var scope: CoroutineScope? = null
|
||||
|
||||
@Volatile
|
||||
private var roundJob: Job? = null
|
||||
private var doneJob: Job? = null
|
||||
|
||||
// Backoff retry after a no-progress, not-exhausted round (relays failed to answer cleanly rather
|
||||
// than empty-EOSE'ing), so a transient connect-storm failure recovers instead of stalling forever.
|
||||
// The user whose window is in flight, read by the done collector when it settles.
|
||||
@Volatile
|
||||
private var retryJob: Job? = null
|
||||
private var windowUser: User? = null
|
||||
|
||||
// Whether a paging window is currently running. Tracked ourselves rather than read from
|
||||
// windowLoad.loading (which starts `true` before any window exists), so the first loadMore actually
|
||||
// starts the window instead of mistaking the construction-time `true` for an in-flight one.
|
||||
@Volatile
|
||||
private var lastRoundUser: User? = null
|
||||
private var windowActive = false
|
||||
|
||||
// The history floor (live-tail boundary) pinned for the current window. startUntil() is `now − 1w`,
|
||||
// which drifts forward in real time — if it were recomputed per assembly, an un-advanced relay's
|
||||
// filter (until = floor) would change every time ANY relay's EOSE triggers invalidateFilters,
|
||||
// re-REQing relays that haven't moved. Pinning it per window keeps those filters stable so only a
|
||||
// relay whose cursor genuinely advanced is re-REQed.
|
||||
@Volatile
|
||||
private var autoLoadAll = false
|
||||
private var windowFloor = 0L
|
||||
|
||||
private fun startUntil() = TimeUtils.now() - AccountGiftWrapsEoseManager.LIVE_TAIL_SECONDS
|
||||
|
||||
private fun floor() = windowFloor.takeIf { it != 0L } ?: startUntil()
|
||||
|
||||
override fun user(key: ChatroomListState) = key.account.userProfile()
|
||||
|
||||
override fun updateFilter(
|
||||
@@ -115,104 +133,80 @@ class ChatroomListNip04HistorySubAssembler(
|
||||
since: SincePerRelayMap?,
|
||||
): List<RelayBasedFilter>? {
|
||||
val user = user(key)
|
||||
if (!key.account.isWriteable() || user.pubkeyHex !in started) {
|
||||
windowLoad.setExpectedRelays(emptySet())
|
||||
return emptyList()
|
||||
}
|
||||
if (!key.account.isWriteable() || user.pubkeyHex !in started) return emptyList()
|
||||
|
||||
// Every relay that still has older history to ask for, each at its own cursor. A relay whose
|
||||
// cursor advanced since the last assembly re-REQs its next page; one still mid-page keeps its open
|
||||
// REQ; a done relay drops out (its REQ closes). This is what lets relays run independently.
|
||||
val homeRelays = key.account.homeRelays.flow.value
|
||||
val dmRelays = key.account.dmRelays.flow.value
|
||||
val active = pager.activeRelays(user.pubkeyHex, (homeRelays + dmRelays).toSet()).toSet()
|
||||
askedRelays[user.pubkeyHex] = active
|
||||
windowLoad.setExpectedRelays(active)
|
||||
if (active.isEmpty()) return emptyList()
|
||||
DmRelayLog.log("rooms.nip04.history", key.account)
|
||||
Log.d("DMPagination") { "[rooms.nip04.history] REQ ${active.size} relay(s), limit=$PAGE_LIMIT fromMe(outbox)=${homeRelays.filter { it in active }.map { it.url }} toMe(inbox)=${dmRelays.filter { it in active }.map { it.url }}" }
|
||||
return homeRelays.filter { it in active }.map {
|
||||
filterNip04DMsFromMe(user, it, since = null, until = pager.untilFor(user.pubkeyHex, it, startUntil()), limit = PAGE_LIMIT)
|
||||
filterNip04DMsFromMe(user, it, since = null, until = pager.untilFor(user.pubkeyHex, it, floor()), limit = PAGE_LIMIT)
|
||||
} +
|
||||
dmRelays.filter { it in active }.map {
|
||||
filterNip04DMsToMe(user, it, since = null, until = pager.untilFor(user.pubkeyHex, it, startUntil()), limit = PAGE_LIMIT)
|
||||
filterNip04DMsToMe(user, it, since = null, until = pager.untilFor(user.pubkeyHex, it, floor()), limit = PAGE_LIMIT)
|
||||
}
|
||||
}
|
||||
|
||||
/** Requests the next backward page from every relay that still has older NIP-04 history. */
|
||||
/** Starts (or resumes) per-relay paging of the NIP-04 history. Idempotent: safe to call again. */
|
||||
fun loadMore(user: User) {
|
||||
if (_exhausted.value) return
|
||||
val account = accounts[user.pubkeyHex] ?: return
|
||||
started.add(user.pubkeyHex)
|
||||
val all = (account.homeRelays.flow.value + account.dmRelays.flow.value).toSet()
|
||||
val active = pager.activeRelays(user.pubkeyHex, all)
|
||||
if (active.isEmpty()) {
|
||||
if (all.isEmpty()) return
|
||||
if (pager.activeRelays(user.pubkeyHex, all).isEmpty()) {
|
||||
// Everything already paged to the bottom.
|
||||
exhaustedByUser[user.pubkeyHex] = true
|
||||
_exhausted.value = true
|
||||
return
|
||||
}
|
||||
val activeSet = active.toSet()
|
||||
if (activeSet == lastAskedActive && lastRoundEventCount == 0) {
|
||||
Log.d("DMPagination") { "[rooms.nip04.history] loadMore skipped — no progress on the same relays" }
|
||||
return
|
||||
}
|
||||
lastAskedActive = activeSet
|
||||
pager.beginRound(user.pubkeyHex, active)
|
||||
lastRoundUser = user
|
||||
_relayCount.value = active.size
|
||||
// Over ALL relays (incl. finished ones) so "reached back to X" stays monotonic and doesn't jump
|
||||
// to a newer date when the deepest relay drops out of the active set.
|
||||
_reachedBack.value = pager.deepestUntil(user.pubkeyHex, all, startUntil())
|
||||
Log.d("DMPagination") { "[rooms.nip04.history] loadMore → ${active.size} active relay(s)" }
|
||||
_exhausted.value = false
|
||||
DmRelayLog.log("rooms.nip04.history", account)
|
||||
windowUser = user
|
||||
scope?.let {
|
||||
ensureRoundCollector(it)
|
||||
windowLoad.startLoading(it)
|
||||
ensureDoneCollector(it)
|
||||
// One window spanning the whole per-relay pagination: it settles a relay only on that relay's
|
||||
// empty-EOSE (done) or when it goes silent/stalled, never on a mid-history page, so the
|
||||
// spinner tracks "is anything still advancing" rather than any single page. Start it only if
|
||||
// none is running — a re-entrant loadMore (the scroll loader re-firing mid-pagination) must
|
||||
// not reset the window and forget the relays that already finished.
|
||||
if (!windowActive) {
|
||||
windowActive = true
|
||||
windowFloor = startUntil()
|
||||
// Populate the relay count BEFORE raising the spinner, so the status card never renders a
|
||||
// "loading from 0 relays" frame between loadingMore flipping true and the first progress.
|
||||
updateStatus(user)
|
||||
_loadingMore.value = true
|
||||
windowLoad.startLoading(it)
|
||||
}
|
||||
windowLoad.setExpectedRelays(all)
|
||||
}
|
||||
updateStatus(user)
|
||||
Log.d("DMPagination") { "[rooms.nip04.history] paging ${all.size} relay(s) independently: ${all.map { it.url }}" }
|
||||
invalidateFilters()
|
||||
}
|
||||
|
||||
/** Pages to the end: each completed round auto-issues the next until exhausted. */
|
||||
fun loadEverything(user: User) {
|
||||
if (_exhausted.value) return
|
||||
autoLoadAll = true
|
||||
loadMore(user)
|
||||
}
|
||||
|
||||
private fun ensureRoundCollector(scope: CoroutineScope) {
|
||||
if (roundJob?.isActive == true) return
|
||||
roundJob =
|
||||
// Mirrors the window's loading state into [_loadingMore] and, when it settles (every relay done or
|
||||
// stalled), flips [exhausted] and clears [windowActive] so the next loadMore can start a fresh window.
|
||||
private fun ensureDoneCollector(scope: CoroutineScope) {
|
||||
if (doneJob?.isActive == true) return
|
||||
doneJob =
|
||||
scope.launch {
|
||||
var wasLoading = false
|
||||
windowLoad.loading.collect { loading ->
|
||||
_loadingMore.value = loading && windowActive
|
||||
if (!loading && wasLoading) {
|
||||
val user = lastRoundUser
|
||||
if (user != null) {
|
||||
val asked = askedRelays[user.pubkeyHex] ?: emptySet()
|
||||
val count = pager.roundEventCount(user.pubkeyHex, asked)
|
||||
lastRoundEventCount = count
|
||||
val account = accounts[user.pubkeyHex]
|
||||
val allRelays = account?.let { (it.homeRelays.flow.value + it.dmRelays.flow.value).toSet() } ?: emptySet()
|
||||
// Exhausted ONLY when every relay returned an empty page + EOSE; CLOSED /
|
||||
// unanswered relays are not finished, so keep loading them.
|
||||
val exhaustedNow = allRelays.isNotEmpty() && pager.activeRelays(user.pubkeyHex, allRelays).isEmpty()
|
||||
exhaustedByUser[user.pubkeyHex] = exhaustedNow
|
||||
_exhausted.value = exhaustedNow
|
||||
// Over ALL relays (incl. finished ones) so the "reached back" date is monotonic.
|
||||
_reachedBack.value = pager.deepestUntil(user.pubkeyHex, allRelays, startUntil())
|
||||
Log.d("DMPagination") { "[rooms.nip04.history] round done: $count event(s), exhausted=$exhaustedNow" }
|
||||
if (autoLoadAll && !exhaustedNow) {
|
||||
loadMore(user)
|
||||
} else if (!exhaustedNow && count == 0) {
|
||||
// No progress and not exhausted: relays failed to answer cleanly rather
|
||||
// than empty-EOSE'ing. Retry after a backoff so a transient failure
|
||||
// recovers, paced so a rate-limited relay isn't hammered.
|
||||
retryJob?.cancel()
|
||||
retryJob =
|
||||
scope.launch {
|
||||
delay(NO_PROGRESS_RETRY_MS)
|
||||
if (!_exhausted.value && !windowLoad.loading.value) {
|
||||
lastAskedActive = emptySet()
|
||||
Log.d("DMPagination") { "[rooms.nip04.history] retry after no-progress round" }
|
||||
loadMore(user)
|
||||
}
|
||||
}
|
||||
}
|
||||
windowActive = false
|
||||
_loadingMore.value = false
|
||||
windowUser?.let { user ->
|
||||
exhaustedByUser[user.pubkeyHex] = true
|
||||
if (activeUser == user.pubkeyHex) _exhausted.value = true
|
||||
updateStatus(user)
|
||||
logSettleSummary(user)
|
||||
}
|
||||
}
|
||||
wasLoading = loading
|
||||
@@ -220,36 +214,71 @@ class ChatroomListNip04HistorySubAssembler(
|
||||
}
|
||||
}
|
||||
|
||||
// WindowLoadTracker reports relays that accepted a REQ then went silent, or never got their REQ out.
|
||||
// We do NOT give up on them (they may simply be slow and need to catch up) — we just record them as
|
||||
// stalled for the logs and let them keep their open subscription.
|
||||
private fun onRelaysStalled(relays: Set<NormalizedRelayUrl>) {
|
||||
started.forEach { pk -> relays.forEach { markStalled(pk, it, "no response (silence/connect timeout)") } }
|
||||
}
|
||||
|
||||
// Records [relay] as not currently advancing for [pk] and logs it once (the first time it stalls in
|
||||
// this window). The relay is kept — it kept its subscription and keeps trying to catch up.
|
||||
private fun markStalled(
|
||||
pk: HexKey,
|
||||
relay: NormalizedRelayUrl,
|
||||
reason: String,
|
||||
) {
|
||||
val firstTime = stalledRelays.getOrPut(pk) { ConcurrentHashMap.newKeySet() }.add(relay)
|
||||
if (firstTime) Log.d("DMPagination") { "[rooms.nip04.history] ${relay.url} stalled — $reason (kept open, still trying)" }
|
||||
}
|
||||
|
||||
private fun updateStatus(user: User) {
|
||||
val account = accounts[user.pubkeyHex]
|
||||
val all = account?.let { (it.homeRelays.flow.value + it.dmRelays.flow.value).toSet() } ?: emptySet()
|
||||
// "Asking N relays" on the status card: the ones still being paged (done relays have dropped out).
|
||||
_relayCount.value = pager.activeRelays(user.pubkeyHex, all).size
|
||||
// Over ALL relays, not just the still-active ones: a relay that finished keeps its deep cursor, so
|
||||
// "reached back to X" stays monotonic instead of jumping back to a newer date when the deepest
|
||||
// relay drops out of the active set.
|
||||
_reachedBack.value = pager.deepestUntil(user.pubkeyHex, all, floor())
|
||||
}
|
||||
|
||||
// A one-line breakdown of where each relay landed when the window settles — the snapshot to reach for
|
||||
// when history didn't load tomorrow: who reached the bottom vs. who is still being retried.
|
||||
private fun logSettleSummary(user: User) {
|
||||
val account = accounts[user.pubkeyHex] ?: return
|
||||
val all = (account.homeRelays.flow.value + account.dmRelays.flow.value).toSet()
|
||||
val done = all.filter { pager.isDone(user.pubkeyHex, it) }.map { it.url }
|
||||
val stillTrying = all.filterNot { pager.isDone(user.pubkeyHex, it) }.map { it.url }
|
||||
Log.d("DMPagination") { "[rooms.nip04.history] settled — done=$done still-trying=$stillTrying" }
|
||||
}
|
||||
|
||||
override fun newSub(key: ChatroomListState): Subscription {
|
||||
val user = user(key)
|
||||
scope = key.account.scope
|
||||
accounts[user.pubkeyHex] = key.account
|
||||
if (activeUser != user.pubkeyHex) {
|
||||
activeUser = user.pubkeyHex
|
||||
// Account switched: repoint the shared display flows to this account's own state.
|
||||
_exhausted.value = exhaustedByUser[user.pubkeyHex] ?: false
|
||||
_relayCount.value = 0
|
||||
_reachedBack.value = null
|
||||
lastAskedActive = emptySet()
|
||||
lastRoundEventCount = -1
|
||||
}
|
||||
return requestNewSubscription(historyListener(user, key))
|
||||
}
|
||||
|
||||
// Flips to exhausted only once every relay has returned an empty page + EOSE. Sets true only.
|
||||
private fun markExhaustedIfAllDone(user: User) {
|
||||
val account = accounts[user.pubkeyHex] ?: return
|
||||
val allRelays = (account.homeRelays.flow.value + account.dmRelays.flow.value).toSet()
|
||||
if (allRelays.isNotEmpty() && pager.activeRelays(user.pubkeyHex, allRelays).isEmpty()) {
|
||||
exhaustedByUser[user.pubkeyHex] = true
|
||||
if (activeUser == user.pubkeyHex) _exhausted.value = true
|
||||
}
|
||||
}
|
||||
|
||||
private fun historyListener(
|
||||
user: User,
|
||||
key: ChatroomListState,
|
||||
): SubscriptionListener =
|
||||
object : SubscriptionListener {
|
||||
override fun onSubscriptionStarted(
|
||||
relay: String,
|
||||
forFilters: List<Filter>,
|
||||
) {
|
||||
windowLoad.onReqSent(relay)
|
||||
}
|
||||
|
||||
override fun onEvent(
|
||||
event: Event,
|
||||
isLive: Boolean,
|
||||
@@ -258,17 +287,27 @@ class ChatroomListNip04HistorySubAssembler(
|
||||
) {
|
||||
windowLoad.onRelayEvent(relay)
|
||||
pager.onEvent(user.pubkeyHex, relay, event.createdAt)
|
||||
stalledRelays[user.pubkeyHex]?.remove(relay)
|
||||
}
|
||||
|
||||
override fun onEose(
|
||||
relay: NormalizedRelayUrl,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
stalledRelays[user.pubkeyHex]?.remove(relay)
|
||||
pager.onEose(user.pubkeyHex, relay)
|
||||
windowLoad.onRelaySettled(relay)
|
||||
if (pager.isDone(user.pubkeyHex, relay)) {
|
||||
// Reached the bottom on this relay: settle it for the spinner, nothing more to ask.
|
||||
windowLoad.onRelaySettled(relay)
|
||||
Log.d("DMPagination") { "[rooms.nip04.history] ${relay.url} reached the bottom (done)" }
|
||||
} else {
|
||||
// This page had events: reset only this relay's tally and let it continue to its next
|
||||
// page immediately, independent of every other relay.
|
||||
pager.beginRound(user.pubkeyHex, listOf(relay))
|
||||
}
|
||||
newEose(key, relay, TimeUtils.now(), forFilters)
|
||||
lastRoundEventCount = -1
|
||||
markExhaustedIfAllDone(user)
|
||||
updateStatus(user)
|
||||
invalidateFilters()
|
||||
}
|
||||
|
||||
override fun onClosed(
|
||||
@@ -276,8 +315,11 @@ class ChatroomListNip04HistorySubAssembler(
|
||||
relay: NormalizedRelayUrl,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
// A relay may demand auth we can't satisfy and CLOSE. It's stalled, not done — keep its
|
||||
// subscription so the pool can re-auth and it can catch up — but don't let it hold the
|
||||
// spinner.
|
||||
windowLoad.onRelaySettled(relay)
|
||||
if (pager.onClosed(user.pubkeyHex, relay)) markExhaustedIfAllDone(user)
|
||||
markStalled(user.pubkeyHex, relay, "CLOSED: $message")
|
||||
}
|
||||
|
||||
override fun onCannotConnect(
|
||||
@@ -285,18 +327,12 @@ class ChatroomListNip04HistorySubAssembler(
|
||||
message: String,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
// Count cannot-connect toward give-up like a CLOSE, so an unreachable relay doesn't keep
|
||||
// exhaustion false forever (otherwise a cold, empty feed retries it endlessly and stays on
|
||||
// the spinner). Once given up, exhaustion completes and the screen resolves.
|
||||
windowLoad.onRelaySettled(relay)
|
||||
if (pager.onClosed(user.pubkeyHex, relay)) markExhaustedIfAllDone(user)
|
||||
markStalled(user.pubkeyHex, relay, "cannot connect: $message")
|
||||
}
|
||||
}
|
||||
|
||||
companion object {
|
||||
private const val PAGE_LIMIT = 10000
|
||||
|
||||
// Backoff before retrying a no-progress, not-exhausted round (transient relay failure).
|
||||
private const val NO_PROGRESS_RETRY_MS = 5_000L
|
||||
}
|
||||
}
|
||||
|
||||
-69
@@ -1,69 +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.amethyst.service.relayClient.eoseManagers
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import org.junit.Assert.assertEquals
|
||||
import org.junit.Assert.assertFalse
|
||||
import org.junit.Assert.assertTrue
|
||||
import org.junit.Test
|
||||
|
||||
class UntilLimitPagerGiveUpTest {
|
||||
private val mine = NormalizedRelayUrl("wss://vitor.nostr1.com/")
|
||||
private val silent = NormalizedRelayUrl("wss://relay.ditto.pub/")
|
||||
private val all = listOf(mine, silent)
|
||||
|
||||
@Test
|
||||
fun givenUpRelayLeavesTheActiveSet() {
|
||||
val pager = UntilLimitPager<String>()
|
||||
|
||||
assertEquals(all, pager.activeRelays("k", all))
|
||||
|
||||
assertTrue("first give-up takes effect", pager.giveUp("k", silent))
|
||||
assertEquals(listOf(mine), pager.activeRelays("k", all))
|
||||
|
||||
assertFalse("giving up twice is a no-op", pager.giveUp("k", silent))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun givingUpEveryRelayExhaustsTheKey() {
|
||||
val pager = UntilLimitPager<String>()
|
||||
|
||||
// mine pages to empty cleanly; the silent relay never answers and is given up.
|
||||
pager.beginRound("k", all)
|
||||
pager.onEose("k", mine) // empty page + EOSE => done
|
||||
pager.giveUp("k", silent)
|
||||
|
||||
assertTrue("no relay left to ask", pager.activeRelays("k", all).isEmpty())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun aRelayThatAlreadyFinishedIsNotMarkedGivenUp() {
|
||||
val pager = UntilLimitPager<String>()
|
||||
|
||||
pager.beginRound("k", listOf(mine))
|
||||
pager.onEose("k", mine) // done
|
||||
|
||||
// A late silence sweep must not "give up" a relay that already finished cleanly.
|
||||
assertFalse(pager.giveUp("k", mine))
|
||||
assertTrue(pager.isDone("k", mine))
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user