refactor(dm): page NIP-17 gift-wrap history per relay, like NIP-04

The gift-wrap history loader was round-based: loadMore asked every active
relay together and the next page only went out after the whole round
settled, so one slow or auth-walled relay throttled the cadence and fast
relays idled until the laggards finished. The per-conversation NIP-04
loader already pages each relay independently; this brings NIP-17 to the
same model.

Now a single loadMore opens one window spanning the whole walk, and each
relay continues itself the instant it EOSEs a non-empty page
(pager.beginRound([relay]) + invalidateFilters, which the sub layer diffs
so only the advanced relay re-REQs). Fast relays race to the bottom while
slow ones catch up in the background. The WindowLoadTracker switches to
tracksReqSends=true with an onAbandoned handler that marks silent/
unreachable relays stalled (kept open, still trying) rather than giving up
on them - exhaustion comes from the window settling via the tracker's
silence + connect-grace backstops, which also removes the need for the old
no-progress retry loop.

Drop the now-dead loadEverything/autoLoadAll (no callers; with independent
paging one loadMore already walks each relay to its bottom). Pin the
history floor per window to keep un-advanced relays' filters stable across
the per-EOSE invalidateFilters, and keep reachedBack monotonic over all
relays.
This commit is contained in:
Claude
2026-06-04 18:36:07 +00:00
parent 89f55a2509
commit 60b8629a27
@@ -41,7 +41,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
@@ -50,15 +49,22 @@ import java.util.concurrent.ConcurrentHashMap
/**
* Loads the account's NIP-17 gift-wrap **history** — everything older than the one-week live tail
* ([AccountGiftWrapsEoseManager]) — by **`until`+`limit` paging, per relay**.
* ([AccountGiftWrapsEoseManager]) — by **`until`+`limit` paging, per relay, independently**.
*
* Idle until a screen calls [loadMore]. Each round asks every not-yet-empty relay for [PAGE_LIMIT]
* gift wraps older than its own cursor (no `since`, so gaps are skipped). When a relay answers an
* empty page with EOSE it is done; otherwise its cursor advances below the oldest wrap it sent. The
* limit is **not** trusted as a stop signal (a relay may cap results on its own) — only an empty page
* is. The whole history is [exhausted] once a full round advances no relay at all (every relay
* empty-EOSE'd or only answered CLOSED), which is the gap-proof "nothing more is reachable" signal the
* old time-slice model couldn't produce.
* There are no lock-step rounds: 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). Fast
* relays race to the bottom of the history in back-to-back pages while slow / auth-walled relays catch
* up at their own pace in the background — none holds the others back. This mirrors the per-conversation
* NIP-04 loader ([com.vitorpamplona.amethyst.ui.screen.loggedIn.chats.privateDM.datasource.ChatroomNip04HistorySubAssembler]).
*
* A relay is *done* once it answers an empty page (nothing older). A relay that won't answer (auth
* CLOSE, unreachable, silent) is marked *stalled* for the logs but keeps its subscription open and keeps
* trying. The [loadingMore] spinner reflects whether anything is still actively advancing across one
* window spanning the whole walk; it clears — and [exhausted] flips — once every relay is either done or
* stalled (the [WindowLoadTracker] settles silent / unreachable relays via its silence + connect-grace
* backstops), without waiting on the slow ones beyond that.
*/
class AccountGiftWrapsHistoryEoseManager(
client: INostrClient,
@@ -66,75 +72,79 @@ class AccountGiftWrapsHistoryEoseManager(
) : PerUserEoseManager<AccountQueryState>(client, allKeys) {
override fun user(key: AccountQueryState) = key.account.userProfile()
// Per-relay cursors, keyed by account pubkey so switching accounts preserves each one's progress.
private val pager = UntilLimitPager<HexKey>()
// Users that have requested history at least once (else the manager stays idle, issuing no REQ).
private val started = ConcurrentHashMap.newKeySet<HexKey>()
// The relays the in-flight round asked, per user — used to tally the round on completion.
private val askedRelays = ConcurrentHashMap<HexKey, Set<NormalizedRelayUrl>>()
// The account behind each user pubkey, captured on subscribe so [loadMore] (UI thread) can read
// the DM relay list without the key.
// The account behind each user pubkey, captured on subscribe so [loadMore] (UI thread) can read the
// DM relay list without the key.
private val accounts = ConcurrentHashMap<HexKey, Account>()
// 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>>()
// This manager is shared across logged-in accounts (one singleton coordinator), so the single
// display flows below must follow whichever account is currently active. Per-account paging
// cursors live in [pager], so switching away and back preserves progress; [exhaustedByUser] lets
// the display flow repoint accurately on switch instead of leaking the previous account's state.
// display flows below must follow whichever account is currently active. Per-account paging cursors
// live in [pager], so switching away and back preserves progress; [exhaustedByUser] lets the display
// flow repoint accurately on switch instead of leaking the previous account's state.
@Volatile
private var activeUser: HexKey? = null
private val exhaustedByUser = ConcurrentHashMap<HexKey, Boolean>()
// No-progress guard: the relay set the last round asked and how many events it returned. If a fresh
// loadMore would ask the same relays and the last round brought nothing (all CLOSED / unanswered),
// we skip it rather than busy-retry — the pool re-auths on the open subscription and its EOSE will
// advance/finish them. Cleared by onEose (any EOSE means something changed).
@Volatile
private var lastAskedActive: Set<NormalizedRelayUrl> = emptySet()
private val windowLoad = WindowLoadTracker("giftwrap.history", tracksReqSends = true, onAbandoned = ::onRelaysStalled)
@Volatile
private var lastRoundEventCount = -1
// 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 windowLoad = WindowLoadTracker("giftwrap.history")
val loadingMore: StateFlow<Boolean> = windowLoad.loading
// True once a full round advanced no relay — nothing older is reachable.
// True once the window settles (every relay done or stalled) — nothing more is reachable right now.
private val _exhausted = MutableStateFlow(false)
val exhausted: StateFlow<Boolean> = _exhausted.asStateFlow()
// Status surfaced to the loading card: how many relays the current page is asking, and the oldest
// point paging has reached (epoch seconds, the deepest cursor).
// Status surfaced to the loading card: how many relays are still being paged, and the oldest point
// paging has reached (epoch seconds, the deepest cursor).
private val _relayCount = MutableStateFlow(0)
val relayCount: StateFlow<Int> = _relayCount.asStateFlow()
private val _reachedBack = MutableStateFlow<Long?>(null)
val reachedBack: StateFlow<Long?> = _reachedBack.asStateFlow()
// Account scope for the watchdog / round collector. Volatile: written on IO (newSub), read on UI.
// Account scope for the done collector. Volatile: written on IO (newSub), read on UI.
@Volatile
private var scope: CoroutineScope? = null
@Volatile
private var roundJob: Job? = null
private var doneJob: Job? = null
// Backoff retry after a round that made no progress but isn't exhausted (relays failed to answer
// cleanly — cannot-connect / CLOSE during a connect storm — rather than empty-EOSE'ing). Without it
// the no-progress guard would never re-fire and a cold, empty feed would stay on the spinner 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
// The user whose round is in flight, read by the round collector on completion.
// 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
// "Load entire history" mode: keep paging to the end without waiting for more scrolling.
// 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
// History starts just below the live tail's one-week floor and pages backward from there.
private fun startUntil() = TimeUtils.now() - AccountGiftWrapsEoseManager.LIVE_TAIL_SECONDS
private fun floor() = windowFloor.takeIf { it != 0L } ?: startUntil()
private fun daysAgo(epochSeconds: Long) = (TimeUtils.now() - epochSeconds) / TimeUtils.ONE_DAY
override fun updateFilter(
@@ -142,115 +152,82 @@ class AccountGiftWrapsHistoryEoseManager(
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 relays = key.account.dmRelays.flow.value
val active = pager.activeRelays(user.pubkeyHex, relays).toSet()
askedRelays[user.pubkeyHex] = active
windowLoad.setExpectedRelays(active)
if (active.isEmpty()) return emptyList()
DmRelayLog.log("giftwrap.history", key.account)
Log.d(TAG) { "[giftwrap.history] REQ ${active.size} relay(s) ${active.map { it.url }}, limit=$PAGE_LIMIT (until ${daysAgo(pager.untilFor(user.pubkeyHex, active.first(), startUntil()))}d…)" }
Log.d(TAG) { "[giftwrap.history] REQ ${active.size} relay(s) ${active.map { it.url }}, limit=$PAGE_LIMIT (until ${daysAgo(pager.untilFor(user.pubkeyHex, active.first(), floor()))}d…)" }
return active.flatMap { relay ->
filterGiftWrapsToPubkey(
relay = relay,
pubkey = user.pubkeyHex,
since = null,
until = pager.untilFor(user.pubkeyHex, relay, startUntil()),
until = pager.untilFor(user.pubkeyHex, relay, floor()),
limit = PAGE_LIMIT,
)
}
}
/** Requests the next backward page from every relay that still has older history. No-op if exhausted. */
/** Starts (or resumes) per-relay paging of the gift-wrap history. Idempotent: safe to call again. */
fun loadMore(user: User) {
if (_exhausted.value) {
Log.d(TAG) { "[giftwrap.history] loadMore ignored — exhausted" }
return
}
val account = accounts[user.pubkeyHex] ?: return
started.add(user.pubkeyHex)
val active = pager.activeRelays(user.pubkeyHex, account.dmRelays.flow.value)
if (active.isEmpty()) {
val allRelays = account.dmRelays.flow.value
if (allRelays.isEmpty()) return
if (pager.activeRelays(user.pubkeyHex, allRelays).isEmpty()) {
// Everything already paged to the bottom.
exhaustedByUser[user.pubkeyHex] = true
_exhausted.value = true
return
}
val activeSet = active.toSet()
if (activeSet == lastAskedActive && lastRoundEventCount == 0) {
// The same relays just returned nothing (e.g. all CLOSED, auth pending). The pool retries
// auth on the open subscription and its EOSE finishes them (see markExhaustedIfAllDone), so
// don't hammer with an identical round. onEose clears this gate when anything changes.
Log.d(TAG) { "[giftwrap.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, 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, account.dmRelays.flow.value, startUntil())
Log.d(TAG) { "[giftwrap.history] loadMore → ${active.size} active relay(s)" }
_exhausted.value = false
DmRelayLog.log("giftwrap.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(allRelays.toSet())
}
updateStatus(user)
Log.d(TAG) { "[giftwrap.history] paging ${allRelays.size} relay(s) independently: ${allRelays.map { it.url }}" }
invalidateFilters()
}
/** Pages to the very end: each completed round auto-issues the next until the history is exhausted. */
fun loadEverything(user: User) {
if (_exhausted.value) return
autoLoadAll = true
Log.d(TAG) { "[giftwrap.history] loadEverything — paging to the end" }
loadMore(user)
}
// Emits the round tally and the exhausted decision when the in-flight load settles. Exhausted ONLY
// when every relay has returned an empty page + EOSE (pager.done): a relay that merely CLOSED (e.g.
// auth-required, before its post-auth retry) or never answered is NOT finished, so we don't call it
// "all caught up" — we keep loading it. In load-all mode, keep paging until that's true.
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 allRelays = accounts[user.pubkeyHex]?.dmRelays?.flow?.value ?: emptySet()
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(TAG) { "[giftwrap.history] round done: $count event(s), exhausted=$exhaustedNow" }
if (autoLoadAll && !exhaustedNow) {
loadMore(user)
} else if (!exhaustedNow && count == 0) {
// No progress and not exhausted: the relays failed to answer cleanly
// (cannot-connect / CLOSE) rather than empty-EOSE'ing. Retry after a
// backoff so a transient failure recovers — paced so a fast-CLOSE
// (rate-limited) relay isn't hammered. Stops once exhausted or loading.
retryJob?.cancel()
retryJob =
scope.launch {
delay(NO_PROGRESS_RETRY_MS)
if (!_exhausted.value && !windowLoad.loading.value) {
lastAskedActive = emptySet()
Log.d(TAG) { "[giftwrap.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
@@ -258,14 +235,41 @@ class AccountGiftWrapsHistoryEoseManager(
}
}
// Flips to exhausted only once every relay has returned an empty page + EOSE (all done). Sets true
// only — the false transitions belong to loadMore / the round collector. Safe off the round path.
private fun markExhaustedIfAllDone(user: User) {
val allRelays = accounts[user.pubkeyHex]?.dmRelays?.flow?.value ?: return
if (allRelays.isNotEmpty() && pager.activeRelays(user.pubkeyHex, allRelays).isEmpty()) {
exhaustedByUser[user.pubkeyHex] = true
if (activeUser == user.pubkeyHex) _exhausted.value = true
}
// 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(TAG) { "[giftwrap.history] ${relay.url} stalled — $reason (kept open, still trying)" }
}
private fun updateStatus(user: User) {
val relays = accounts[user.pubkeyHex]?.dmRelays?.flow?.value ?: emptySet()
// "Asking N relays" on the status card: the ones still being paged (done relays have dropped out).
_relayCount.value = pager.activeRelays(user.pubkeyHex, relays).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, relays, 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 relays = accounts[user.pubkeyHex]?.dmRelays?.flow?.value ?: return
val done = relays.filter { pager.isDone(user.pubkeyHex, it) }.map { it.url }
val stillTrying = relays.filterNot { pager.isDone(user.pubkeyHex, it) }.map { it.url }
Log.d(TAG) { "[giftwrap.history] settled — done=$done still-trying=$stillTrying" }
}
override fun newSub(key: AccountQueryState): Subscription {
@@ -278,8 +282,6 @@ class AccountGiftWrapsHistoryEoseManager(
_exhausted.value = exhaustedByUser[user.pubkeyHex] ?: false
_relayCount.value = 0
_reachedBack.value = null
lastAskedActive = emptySet()
lastRoundEventCount = -1
}
return requestNewSubscription(historyListener(user, key))
}
@@ -289,6 +291,13 @@ class AccountGiftWrapsHistoryEoseManager(
key: AccountQueryState,
): SubscriptionListener =
object : SubscriptionListener {
override fun onSubscriptionStarted(
relay: String,
forFilters: List<Filter>,
) {
windowLoad.onReqSent(relay)
}
override fun onEvent(
event: Event,
isLive: Boolean,
@@ -297,21 +306,27 @@ class AccountGiftWrapsHistoryEoseManager(
) {
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(TAG) { "[giftwrap.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)
// An EOSE means this relay changed (finished, or delivered a page) — clear the
// no-progress gate so the next loadMore can continue, even off the round path.
lastRoundEventCount = -1
// A post-auth empty EOSE can land after the round already settled on the earlier CLOSED;
// flip to exhausted the moment this finishes the last relay, not only at round end.
markExhaustedIfAllDone(user)
updateStatus(user)
invalidateFilters()
}
override fun onClosed(
@@ -319,12 +334,11 @@ class AccountGiftWrapsHistoryEoseManager(
relay: NormalizedRelayUrl,
forFilters: List<Filter>?,
) {
// CLOSED (e.g. auth-required) is not "empty": don't mark the relay done — it may answer
// after the auth handshake. It just settles the load so the spinner can clear. But if it
// keeps rejecting us (auth we can't satisfy), the pager eventually gives up on it; once
// that's the last blocker, exhaustion can complete.
// A relay (e.g. an author's) 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(
@@ -332,12 +346,8 @@ class AccountGiftWrapsHistoryEoseManager(
message: String,
forFilters: List<Filter>?,
) {
// Cannot-connect is also "no answer": count it toward give-up like a CLOSE, so a relay
// that's unreachable (down / blocked) doesn't keep exhaustion false forever — otherwise a
// cold, empty feed would retry it endlessly and stay on the spinner. Once it's given up,
// exhaustion can complete and the screen resolves (to the loaded rooms, or empty + retry).
windowLoad.onRelaySettled(relay)
if (pager.onClosed(user.pubkeyHex, relay)) markExhaustedIfAllDone(user)
markStalled(user.pubkeyHex, relay, "cannot connect: $message")
}
}
@@ -348,8 +358,5 @@ class AccountGiftWrapsHistoryEoseManager(
// relay allows it. A relay returning fewer is treated as its own cap, NOT as "nothing more" —
// only an empty page + EOSE ends a relay.
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
}
}