feat: page each NIP-04 DM relay independently to converge on one window

Replaces the lock-step round model (every relay advanced one page per
global round, gated by the slowest) with per-relay continuous paging:
each relay continues to its next page the instant it EOSEs, off its own
cursor. The subscription layer diffs per relay, so re-issuing only
re-REQs the relay whose cursor moved — others' in-flight REQs are
untouched. Fast relays race to the bottom in back-to-back pages while
slow / auth-walled relays catch up at their own pace; none are
abandoned (this reverts the give-up behaviour — slow relays keep their
subscription open and keep trying), so every relay converges on the same
window.

A relay is done on an empty page; one that won't answer (auth CLOSE,
unreachable, silent) is marked stalled but keeps trying. loadingMore
clears once every relay is done or stalled. Exposes per-relay
RelayPagingProgress (reached-back / done / stalled) for the upcoming
in-stream progress markers.

Splits the convo widen loop so each protocol pages on its own loader
state — NIP-04's continuous loading no longer starves gift-wrap paging.

https://claude.ai/code/session_01B1fmmmX8JjQWH3amMLdvcW
This commit is contained in:
Claude
2026-06-03 00:29:53 +00:00
parent 07b0c87c38
commit cbee966bb7
2 changed files with 136 additions and 144 deletions
@@ -161,8 +161,10 @@ private const val PREFETCH_OLDER_MESSAGES = 3
* room, so this advances the shared account-wide history window and the conversation's messages
* surface as its pages are decrypted.
*
* Both protocols advance via their history managers' `loadMore`; the step is gated on BOTH loaders
* being idle, so it never outruns the slower one, and stops once both report exhausted.
* Each protocol advances via its own history manager's `loadMore`, gated only on ITS OWN loader/
* exhausted state — so a slow protocol (e.g. NIP-04 waiting on a sluggish correspondent relay) never
* holds back the other. NIP-04 then pages every relay independently to completion on its own; gift
* wraps stay round-driven and re-step here as the oldest end stays in view.
*/
@Composable
private fun LoadOlderMessagesWhenScrolling(
@@ -173,7 +175,7 @@ private fun LoadOlderMessagesWhenScrolling(
val nip04History = remember(accountViewModel) { accountViewModel.dataSources().chatroom.nip04History }
LaunchedEffect(listState, giftWrapsHistory, nip04History) {
combine(
val wantMore =
snapshotFlow {
val info = listState.layoutInfo
val total = info.totalItemsCount
@@ -181,21 +183,24 @@ private fun LoadOlderMessagesWhenScrolling(
// The oldest end is in view (no overflow requirement, so a one-message thread that
// can't scroll still qualifies and walks history to its start).
total > 0 && lastVisible >= total - PREFETCH_OLDER_MESSAGES
},
giftWrapsHistory.loadingMore,
nip04History.loadingMore,
giftWrapsHistory.exhausted,
nip04History.exhausted,
) { wantMore, loadingGiftWraps, loadingNip04, giftWrapsExhausted, nip04Exhausted ->
// Keep paging while either protocol still has older history to reach.
wantMore && !loadingGiftWraps && !loadingNip04 && !(giftWrapsExhausted && nip04Exhausted)
}.distinctUntilChanged()
.filter { it }
.collect {
Log.d("DMPagination") { "convo: widen (oldest in view) → loadMore" }
}.distinctUntilChanged()
launch {
combine(wantMore, giftWrapsHistory.loadingMore, giftWrapsHistory.exhausted) { want, loading, exhausted ->
want && !loading && !exhausted
}.distinctUntilChanged().filter { it }.collect {
Log.d("DMPagination") { "convo: widen (oldest in view) → giftwrap loadMore" }
giftWrapsHistory.loadMore(accountViewModel.userProfile())
}
}
launch {
combine(wantMore, nip04History.loadingMore, nip04History.exhausted) { want, loading, exhausted ->
want && !loading && !exhausted
}.distinctUntilChanged().filter { it }.collect {
Log.d("DMPagination") { "convo: widen (oldest in view) → nip04 loadMore" }
nip04History.loadMore()
}
}
}
}
@@ -20,7 +20,6 @@
*/
package com.vitorpamplona.amethyst.ui.screen.loggedIn.chats.privateDM.datasource
import com.vitorpamplona.amethyst.service.relayClient.eoseManagers.DmRelayLog
import com.vitorpamplona.amethyst.service.relayClient.eoseManagers.PerUserAndFollowListEoseManager
import com.vitorpamplona.amethyst.service.relayClient.eoseManagers.UntilLimitPager
import com.vitorpamplona.amethyst.service.relayClient.eoseManagers.WindowLoadTracker
@@ -45,11 +44,30 @@ import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.launch
import java.util.concurrent.ConcurrentHashMap
/** How far back one relay has paged a conversation, for the per-relay progress markers. */
data class RelayPagingProgress(
// The oldest createdAt this relay has loaded down to (its `until` cursor). The marker sits here and
// slides down (older) as the relay pages further back.
val reachedUntil: Long,
// The relay answered an empty page: it has nothing older, it has reached the bottom of its window.
val done: Boolean,
// The relay isn't answering right now (auth-walled CLOSE / unreachable / slow). It is NOT abandoned
// — its subscription stays open and it keeps trying to catch up — but it isn't currently advancing.
val stalled: Boolean,
)
/**
* Loads older NIP-04 DMs (kind 4) for one conversation by `until`+`limit` paging, per relay, scoped to
* the two participants. Same gap-proof model as the gift-wrap history (a relay is done on an empty
* page + EOSE; [exhausted] once a round advances no relay), keyed per conversation. Idle until
* [loadMore].
* Loads older NIP-04 DMs (kind 4) for one conversation by `until`+`limit` paging — **per relay,
* independently**. There are no lock-step rounds: every relay drives its own pages off its own cursor,
* continuing the instant it EOSEs (the subscription layer diffs per relay, so re-issuing only re-REQs
* the relay whose cursor moved; the others' in-flight REQs are untouched). Fast relays race to the
* bottom of the conversation in a few back-to-back pages while slow / auth-walled relays catch up at
* their own pace in the background — none are abandoned, so they all converge on the same window.
*
* 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 markers but keeps its subscription open and
* keeps trying. The [loadingMore] spinner reflects whether anything is still actively advancing; it
* clears once every relay is either done or stalled, without waiting on the slow ones beyond that.
*/
class ChatroomNip04HistorySubAssembler(
client: INostrClient,
@@ -67,9 +85,12 @@ class ChatroomNip04HistorySubAssembler(
private val pager = UntilLimitPager<ConvoKey>()
private val started = ConcurrentHashMap.newKeySet<ConvoKey>()
private val askedRelays = ConcurrentHashMap<ConvoKey, Set<NormalizedRelayUrl>>()
private val windowLoad = WindowLoadTracker("convo.nip04.history", tracksReqSends = true, onAbandoned = ::onRelaysAbandoned)
// Relays currently not advancing for a conversation (auth CLOSE / unreachable / silent). Tracked for
// the progress markers; these relays are NOT given up — they keep their subscription and keep trying.
private val stalledRelays = ConcurrentHashMap<ConvoKey, MutableSet<NormalizedRelayUrl>>()
private val windowLoad = WindowLoadTracker("convo.nip04.history", tracksReqSends = true, onAbandoned = ::onRelaysStalled)
val loadingMore: StateFlow<Boolean> = windowLoad.loading
private val _exhausted = MutableStateFlow(false)
@@ -81,28 +102,21 @@ class ChatroomNip04HistorySubAssembler(
private val _reachedBack = MutableStateFlow<Long?>(null)
val reachedBack: StateFlow<Long?> = _reachedBack.asStateFlow()
// Per-relay paging progress for the conversation on screen — the data the in-stream markers render.
private val _relayProgress = MutableStateFlow<Map<NormalizedRelayUrl, RelayPagingProgress>>(emptyMap())
val relayProgress: StateFlow<Map<NormalizedRelayUrl, RelayPagingProgress>> = _relayProgress.asStateFlow()
// Shared across accounts/conversations (singleton coordinator): repoint the display flows to the
// conversation now on screen instead of leaking the previous one's state. Cursors live in [pager].
@Volatile
private var activeConvo: ConvoKey? = null
private val exhaustedByConvo = ConcurrentHashMap<ConvoKey, Boolean>()
// No-progress guard (see the gift-wrap history manager's twin): skip re-issuing an identical round
// that brought nothing; cleared by onEose.
@Volatile
private var lastAskedActive: Set<NormalizedRelayUrl> = emptySet()
@Volatile
private var lastRoundEventCount = -1
@Volatile
private var scope: CoroutineScope? = null
@Volatile
private var roundJob: Job? = null
@Volatile
private var autoLoadAll = false
private var doneJob: Job? = null
private fun startUntil() = TimeUtils.now() - AccountGiftWrapsEoseManager.LIVE_TAIL_SECONDS
@@ -116,115 +130,104 @@ class ChatroomNip04HistorySubAssembler(
): List<RelayBasedFilter>? {
val pk = convoKey(key)
val relays = nip04DMRelays(key.room.users, key.account)
if (!key.account.isWriteable() || pk !in started || relays == null) {
windowLoad.setExpectedRelays(emptySet())
return emptyList()
}
if (!key.account.isWriteable() || pk !in started || relays == null) 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 active = pager.activeRelays(pk, relays.all).toSet()
askedRelays[pk] = active
windowLoad.setExpectedRelays(active)
if (active.isEmpty()) return emptyList()
DmRelayLog.log("convo.nip04.history", key.account)
Log.d("DMPagination") { "[convo.nip04.history] REQ ${active.size} relay(s), limit=$PAGE_LIMIT fromMe(outbox)=${relays.fromMeRelays.keys.intersect(active).map { it.url }} toMe(inbox)=${relays.toMeRelays.keys.intersect(active).map { it.url }}" }
val activeRelays =
val scoped =
Nip04DmRelays(
toMeRelays = relays.toMeRelays.filterKeys { it in active },
fromMeRelays = relays.fromMeRelays.filterKeys { it in active },
)
return filterNip04DMsHistory(key.account, activeRelays, PAGE_LIMIT) { relay ->
return filterNip04DMsHistory(key.account, scoped, PAGE_LIMIT) { relay ->
pager.untilFor(pk, relay, startUntil())
}
}
/** Requests the next backward page for every open conversation that still has older history. */
/** Starts (or resumes) per-relay paging for every open conversation. Idempotent: safe to call again. */
fun loadMore() {
if (_exhausted.value) return
// Gather the active (not-finished) relays per open conversation first, so the no-progress guard
// is checked before any beginRound (which would otherwise reset round tallies prematurely).
val perKeyActive = mutableListOf<Pair<ConvoKey, List<NormalizedRelayUrl>>>()
val activeUnion = mutableSetOf<NormalizedRelayUrl>()
allKeys().forEach { key ->
val keys = allKeys()
val fullRelays = mutableSetOf<NormalizedRelayUrl>()
var anyActive = false
keys.forEach { key ->
val relays = nip04DMRelays(key.room.users, key.account) ?: return@forEach
val pk = convoKey(key)
started.add(pk)
val active = pager.activeRelays(pk, relays.all)
if (active.isNotEmpty()) {
perKeyActive.add(pk to active)
activeUnion.addAll(active)
}
fullRelays.addAll(relays.all)
if (pager.activeRelays(pk, relays.all).isNotEmpty()) anyActive = true
}
if (perKeyActive.isEmpty()) {
if (fullRelays.isEmpty()) return
if (!anyActive) {
// Everything already paged to the bottom.
activeConvo?.let { exhaustedByConvo[it] = true }
_exhausted.value = true
return
}
if (activeUnion == lastAskedActive && lastRoundEventCount == 0) {
// The same relays just returned nothing two rounds running — they're unreadable: a
// correspondent's auth-walled relay that only ever CLOSEs (ditto: "all authors must be
// authenticated"), or one perpetually stuck reconnecting. Give up on them so the
// conversation reports as finished instead of lingering forever on a "N relays" count
// with no progress bar. (My own reachable relays empty-EOSE to `done` and never land here.)
Log.d("DMPagination") { "[convo.nip04.history] giving up — no progress on the same relays ${activeUnion.map { it.url }}" }
perKeyActive.forEach { (pk, active) -> active.forEach { pager.giveUp(pk, it) } }
markExhaustedIfAllDone()
return
}
lastAskedActive = activeUnion
var totalRelays = 0
var deepest: Long? = null
perKeyActive.forEach { (pk, active) ->
pager.beginRound(pk, active)
totalRelays += active.size
pager.deepestUntil(pk, active, startUntil())?.let { d ->
deepest = deepest?.let { minOf(it, d) } ?: d
}
}
_relayCount.value = totalRelays
_reachedBack.value = deepest
Log.d("DMPagination") { "[convo.nip04.history] loadMore" }
_exhausted.value = false
_relayCount.value = fullRelays.size
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 round.
if (!windowLoad.loading.value) windowLoad.startLoading(it)
windowLoad.setExpectedRelays(fullRelays)
}
publishProgress()
Log.d("DMPagination") { "[convo.nip04.history] paging ${fullRelays.size} relay(s) independently" }
invalidateFilters()
}
/** Pages to the end: each completed round auto-issues the next until exhausted. */
fun loadEverything() {
if (_exhausted.value) return
autoLoadAll = true
loadMore()
}
/** Per-relay paging already runs to completion on its own, so loading everything is just [loadMore]. */
fun loadEverything() = loadMore()
private fun ensureRoundCollector(scope: CoroutineScope) {
if (roundJob?.isActive == true) return
roundJob =
// Flips [exhausted] when the window settles (every relay done or stalled) and back to false when a
// fresh page starts. Tied to the spinner so "nothing is advancing" and "caught up" stay consistent.
private fun ensureDoneCollector(scope: CoroutineScope) {
if (doneJob?.isActive == true) return
doneJob =
scope.launch {
var wasLoading = false
windowLoad.loading.collect { loading ->
if (!loading && wasLoading) {
val count = started.sumOf { pk -> pager.roundEventCount(pk, askedRelays[pk] ?: emptySet()) }
lastRoundEventCount = count
// Exhausted ONLY when every open conversation's relays have all returned an
// empty page + EOSE; a CLOSED / unanswered relay isn't finished, so keep loading.
val keys = allKeys()
val exhaustedNow =
keys.isNotEmpty() &&
keys.none { key ->
val relays = nip04DMRelays(key.room.users, key.account)
relays != null && pager.activeRelays(convoKey(key), relays.all).isNotEmpty()
}
activeConvo?.let { exhaustedByConvo[it] = exhaustedNow }
_exhausted.value = exhaustedNow
_reachedBack.value = started.mapNotNull { pk -> pager.deepestUntil(pk, askedRelays[pk] ?: emptySet(), startUntil()) }.minOrNull()
Log.d("DMPagination") { "[convo.nip04.history] round done: $count event(s), exhausted=$exhaustedNow" }
if (autoLoadAll && !exhaustedNow) loadMore()
activeConvo?.let { exhaustedByConvo[it] = true }
_exhausted.value = true
publishProgress()
Log.d("DMPagination") { "[convo.nip04.history] all relays settled (done or stalled)" }
}
wasLoading = loading
}
}
}
// 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 markers and let them keep their open subscription.
private fun onRelaysStalled(relays: Set<NormalizedRelayUrl>) {
started.forEach { pk -> stalledRelays.getOrPut(pk) { ConcurrentHashMap.newKeySet() }.addAll(relays) }
publishProgress()
}
private fun publishProgress() {
val pk = activeConvo ?: return
val relays = allKeys().firstOrNull { convoKey(it) == pk }?.let { nip04DMRelays(it.room.users, it.account) } ?: return
val stalled = stalledRelays[pk] ?: emptySet()
val start = startUntil()
_relayProgress.value =
relays.all.associateWith { relay ->
RelayPagingProgress(
reachedUntil = pager.untilFor(pk, relay, start),
done = pager.isDone(pk, relay),
stalled = relay in stalled && !pager.isDone(pk, relay),
)
}
_reachedBack.value = pager.deepestUntil(pk, relays.all, start)
}
override fun newSub(key: ChatroomQueryState): Subscription {
scope = key.account.scope
val pk = convoKey(key)
@@ -234,40 +237,11 @@ class ChatroomNip04HistorySubAssembler(
_exhausted.value = exhaustedByConvo[pk] ?: false
_relayCount.value = 0
_reachedBack.value = null
lastAskedActive = emptySet()
lastRoundEventCount = -1
_relayProgress.value = emptyMap()
}
return requestNewSubscription(historyListener(key))
}
// A relay accepted the REQ but never answered (auth-walled / dead): drop it from every open
// conversation's pager so it stops blocking the relay count and exhaustion on the next round. May
// complete exhaustion right away if it was the last relay still holding a thread open.
private fun onRelaysAbandoned(relays: Set<NormalizedRelayUrl>) {
var gaveUp = false
started.forEach { pk ->
val asked = askedRelays[pk] ?: return@forEach
relays.forEach { if (it in asked && pager.giveUp(pk, it)) gaveUp = true }
}
if (gaveUp) markExhaustedIfAllDone()
}
// Flips to exhausted only once every open conversation's relays have all returned an empty page +
// EOSE. Sets true only — false transitions belong to loadMore / the round collector.
private fun markExhaustedIfAllDone() {
val keys = allKeys()
val allDone =
keys.isNotEmpty() &&
keys.none { key ->
val relays = nip04DMRelays(key.room.users, key.account)
relays != null && pager.activeRelays(convoKey(key), relays.all).isNotEmpty()
}
if (allDone) {
activeConvo?.let { exhaustedByConvo[it] = true }
_exhausted.value = true
}
}
private fun historyListener(key: ChatroomQueryState): SubscriptionListener {
val pk = convoKey(key)
return object : SubscriptionListener {
@@ -286,17 +260,26 @@ class ChatroomNip04HistorySubAssembler(
) {
windowLoad.onRelayEvent(relay)
pager.onEvent(pk, relay, event.createdAt)
stalledRelays[pk]?.remove(relay)
}
override fun onEose(
relay: NormalizedRelayUrl,
forFilters: List<Filter>?,
) {
stalledRelays[pk]?.remove(relay)
pager.onEose(pk, relay)
windowLoad.onRelaySettled(relay)
if (pager.isDone(pk, relay)) {
// Reached the bottom on this relay: settle it for the spinner, nothing more to ask.
windowLoad.onRelaySettled(relay)
} 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(pk, listOf(relay))
}
newEose(key, relay, TimeUtils.now(), forFilters)
lastRoundEventCount = -1
markExhaustedIfAllDone()
publishProgress()
invalidateFilters()
}
override fun onClosed(
@@ -304,10 +287,12 @@ class ChatroomNip04HistorySubAssembler(
relay: NormalizedRelayUrl,
forFilters: List<Filter>?,
) {
// A relay (e.g. the correspondent'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)
// A relay (e.g. the correspondent's) may demand auth we can't satisfy and CLOSE every
// round; once the pager gives up on it, it stops blocking this thread's exhaustion.
if (pager.onClosed(pk, relay)) markExhaustedIfAllDone()
stalledRelays.getOrPut(pk) { ConcurrentHashMap.newKeySet() }.add(relay)
publishProgress()
}
override fun onCannotConnect(
@@ -316,6 +301,8 @@ class ChatroomNip04HistorySubAssembler(
forFilters: List<Filter>?,
) {
windowLoad.onRelaySettled(relay)
stalledRelays.getOrPut(pk) { ConcurrentHashMap.newKeySet() }.add(relay)
publishProgress()
}
}
}