fix(nip29): scope group subscriptions per channel so Messages updates live

The Messages list showed a stale last message for every NIP-29 group while the
open chat screen stayed live. The joined-group tail batched every group on a
host relay into ONE filter carrying all their ids in `#h`. That is valid NIP-01
and relays answer it correctly for stored queries — which is what made it so
confusing: the boot backfill populated every group, so the list looked right
until it needed to change.

`block/buzz` indexes each live subscription under a single channel uuid resolved
from `#h`. When two or more distinct ids appear anywhere across a subscription's
filters, `extract_channel_id_from_filters` returns None and the subscription is
registered as *global* — and global subscriptions deliberately never receive
channel-scoped events, guarding against leaking private channel content to a
subscriber whose membership was not checked per channel. So the batched tail
backfilled at EOSE and then went permanently deaf.

Measured on device: same relay, same connection, same subscription id, only the
`#h` count changed — one value delivered live, two delivered nothing.

The resolver scans every filter of a subscription, so splitting into one filter
per group is not enough; each channel needs its own subscription. The EOSE
managers are therefore keyed on GroupId, and the preloads mount one subscription
per joined channel. Relays leave room for this: buzz allows max_subscriptions
1024 against max_filters 10, so this also sidesteps the filter cap for anyone in
more than ten groups.

Also folded into the same per-channel subscriptions:

- Group activity addressed to me (reactions/zaps/replies) moved off the
  account-wide notifications subscription. That one also carries inbox filters
  with no `#h` at all, and buzz forces a subscription global on the first
  channel-less filter, so no reshaping there could ever have worked.
- Reactions and deletions (kinds 5/7/9005) now ride the channel's own `#h`
  subscription, the shape Buzz's own client uses (`channelEventKinds`). Amethyst
  otherwise learned about them only through the shared `#e` EventFinder query,
  which carries no `#h` and is therefore global — so reaction chips appeared
  only on a re-query, never as they happened. Kept out of the timeline kind set
  so reactions cannot consume a history page's limit and walk the `until` cursor
  past undelivered messages.
- A `limit = 1` preview filter with no time floor. The tail floors at now-7d to
  keep recent chat warm; on its own that stranded any channel quiet for longer
  on a "No messages yet" placeholder sorted to the bottom by createdAt 0. Every
  other roster-driven protocol on that screen already bounds by count for this
  reason (NIP-28 `limit = 1`, Concord `limit = 10`) — their row set comes from a
  list event, unlike NIP-17/NIP-04 whose rooms are discovered from the messages
  and so need the backward pagers.

The Messages row now says "Loading" instead of "No messages yet" until the fleet
settles, so an in-flight channel is not mistaken for an empty one. That required
the shared WindowLoadTracker to describe the whole joined fleet rather than
whichever channel rebuilt last.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Vitor Pamplona
2026-07-26 12:44:05 -04:00
co-authored by Claude Opus 5
parent 683d90763b
commit f434e98b00
9 changed files with 372 additions and 159 deletions
@@ -20,8 +20,6 @@
*/
package com.vitorpamplona.amethyst.service.relayClient.reqCommand.account.nip01Notifications
import com.vitorpamplona.amethyst.commons.model.buzz.BuzzDmChannels
import com.vitorpamplona.amethyst.commons.model.buzz.BuzzDmRegistry
import com.vitorpamplona.amethyst.model.User
import com.vitorpamplona.amethyst.service.relayClient.eoseManagers.PerUserEoseManager
import com.vitorpamplona.amethyst.service.relayClient.reqCommand.account.AccountQueryState
@@ -29,7 +27,6 @@ import com.vitorpamplona.amethyst.service.relays.SincePerRelayMap
import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient
import com.vitorpamplona.quartz.nip01Core.relay.client.pool.RelayBasedFilter
import com.vitorpamplona.quartz.nip01Core.relay.client.subscriptions.Subscription
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.FlowPreview
import kotlinx.coroutines.Job
@@ -84,45 +81,17 @@ class AccountNotificationsEoseFromInboxRelaysManager(
)
}
// NIP-29 group activity (reactions/replies to my messages) lives on the group's host relay,
// not my inbox relays — poll each host relay for those, scoped by `#h` to my joined groups.
val groups =
key.account.relayGroupList.liveRelayGroupList.value
.groupBy({ it.relayUrl }, { it.groupId })
.flatMap { (relayUrl, groupIds) ->
val relay = RelayUrlNormalizer.normalizeOrNull(relayUrl) ?: return@flatMap emptyList()
filterGroupNotificationsToPubkey(
relay = relay,
pubkey = user(key).pubkeyHex,
groupIds = groupIds.distinct(),
// Same reasoning as the inbox filters above: `#p` + `#h` + `limit = 200`
// already bound this, so a week floor only hides older group activity.
since = since?.get(relay)?.time ?: pagingBoundary,
)
}
// Buzz DM channels are NOT in the published group list (membership is server-side, tracked in
// BuzzDmChannels), so the joined-group query above skips them. Poll each DM host relay the same
// way — `#p` = me + `#h` = my DM channels — so a reaction/zap/repost on my DM message surfaces
// in notifications, not only as a chip inside the open conversation. Hidden DMs are excluded.
val myPubkey = user(key).pubkeyHex
val hiddenDms = BuzzDmRegistry.hiddenFor(myPubkey)
val dmGroups =
BuzzDmChannels
.channelsFor(myPubkey)
.filterKeys { it !in hiddenDms }
.entries
.groupBy({ it.value }, { it.key })
.flatMap { (relay, channelIds) ->
filterGroupNotificationsToPubkey(
relay = relay,
pubkey = myPubkey,
groupIds = channelIds.distinct(),
since = since?.get(relay)?.time ?: pagingBoundary,
)
}
return inbox + groups + dmGroups
// NIP-29 group activity (reactions/replies to my messages) is deliberately NOT requested here.
// It lives on the group's host relay and used to be one `#h` filter per relay carrying every
// joined group id — but this subscription also carries the inbox filters above, which have no
// `#h` at all, and `block/buzz` downgrades any subscription with a channel-less (or multi-
// channel) filter to "global", a class that by design never receives channel-scoped events. So
// those filters answered the stored query at EOSE and then went deaf until the next launch.
//
// Group and Buzz-DM activity now rides the per-channel subscriptions that are already scoped to
// exactly one channel — RelayGroupJoinedChatTailSubAssembler and BuzzDmJoinedChatTailSubAssembler,
// both mounted app-wide from LoggedInPage, so coverage is unchanged and delivery is live.
return inbox
}
val userJobMap = mutableMapOf<User, List<Job>>()
@@ -138,21 +107,8 @@ class AccountNotificationsEoseFromInboxRelaysManager(
invalidateFilters()
}
},
// Re-subscribe when I join/leave a group so its host relay is added to (or dropped
// from) the notification query.
key.account.scope.launch(Dispatchers.IO) {
key.account.relayGroupList.liveRelayGroupList.sample(1000).collectLatest {
invalidateFilters()
}
},
// Re-subscribe when a Buzz DM is discovered/hidden so its host relay is polled for
// reactions/zaps on my DM messages.
key.account.scope.launch(Dispatchers.IO) {
BuzzDmChannels.flow.sample(1000).collectLatest { invalidateFilters() }
},
key.account.scope.launch(Dispatchers.IO) {
BuzzDmRegistry.hidden.sample(1000).collectLatest { invalidateFilters() }
},
// No group/Buzz-DM watchers here any more: those filters moved to the per-channel
// subscriptions, which mount and unmount with the channel itself.
key.account.scope.launch(Dispatchers.IO) {
key.feedContentStates.notifications.lastNoteCreatedAtWhenFullyLoaded.sample(5000).collectLatest {
invalidateFilters()
@@ -28,18 +28,20 @@ import com.vitorpamplona.amethyst.commons.relayClient.paging.WindowLoadTracker
import com.vitorpamplona.amethyst.commons.relayClient.paging.trackingListener
import com.vitorpamplona.amethyst.model.Account
import com.vitorpamplona.amethyst.service.relayClient.eoseManagers.PerUniqueIdEoseManager
import com.vitorpamplona.amethyst.service.relayClient.reqCommand.account.nip01Notifications.filterGroupNotificationsToPubkey
import com.vitorpamplona.amethyst.service.relays.SincePerRelayMap
import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient
import com.vitorpamplona.quartz.nip01Core.relay.client.pool.RelayBasedFilter
import com.vitorpamplona.quartz.nip01Core.relay.client.subscriptions.Subscription
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.nip51Lists.simpleGroupList.GroupTag
import com.vitorpamplona.quartz.nip29RelayGroups.GroupId
import com.vitorpamplona.quartz.utils.TimeUtils
import kotlinx.coroutines.flow.StateFlow
/** One screen's request to keep the viewer's Buzz DM channels' recent chat live. */
/** One Buzz DM channel whose recent chat is kept live. */
class BuzzDmJoinedChatTailQueryState(
val account: Account,
val groupId: GroupId,
)
/**
@@ -71,7 +73,7 @@ class BuzzDmJoinedChatTailFilterAssembler(
class BuzzDmJoinedChatTailSubAssembler(
client: INostrClient,
allKeys: () -> Set<BuzzDmJoinedChatTailQueryState>,
) : PerUniqueIdEoseManager<BuzzDmJoinedChatTailQueryState, Account>(client, allKeys) {
) : PerUniqueIdEoseManager<BuzzDmJoinedChatTailQueryState, GroupId>(client, allKeys) {
private val windowLoad = WindowLoadTracker("buzzDm.preview.live")
val loadingMore: StateFlow<Boolean> = windowLoad.loading
@@ -79,23 +81,35 @@ class BuzzDmJoinedChatTailSubAssembler(
key: BuzzDmJoinedChatTailQueryState,
since: SincePerRelayMap?,
): List<RelayBasedFilter>? {
// The preload only mounts visible channels; re-checked here so hiding a conversation
// empties its subscription even if the key outlives the recomposition.
val me = key.account.userProfile().pubkeyHex
val hidden = BuzzDmRegistry.hiddenFor(me)
val channels = BuzzDmChannels.channelsFor(me).filterKeys { it !in hidden }
if (channels.isEmpty()) {
windowLoad.setExpectedRelays(emptySet())
return null
}
if (key.groupId.id in BuzzDmRegistry.hiddenFor(me)) return null
// Reuse the joined-group tail builder: one #h filter per host relay carrying every DM channel id
// on it, bounded by the shared recent floor (no per-channel limit, reconnect-safe).
val asTags = channels.map { (channelId, relay) -> GroupTag(channelId, relay.url) }
val filters = buildRelayGroupJoinedChatTailFilters(asTags, DmHistoryTuning.recentBoundary())
windowLoad.setExpectedRelays(filters.mapTo(mutableSetOf()) { it.relay })
return filters
// Same one-subscription-per-channel shape as the joined-group tail: a Buzz DM is a
// relay-authoritative NIP-29 group, so batching ids into one `#h` would register the
// subscription as global on buzz and never deliver live messages. See
// [buildRelayGroupJoinedChatTailFilter]. The second filter carries the DM's activity
// addressed to me (reactions/zaps/reports) — same channel, so the subscription stays scoped.
val relay = key.groupId.relayUrl
windowLoad.setExpectedRelays(setOf(relay))
return listOf(
buildRelayGroupJoinedChatTailFilter(key.groupId, DmHistoryTuning.recentBoundary()),
// Newest message at any age — a DM conversation dormant for over a week must still show its
// last message on Messages, not a placeholder. See [buildRelayGroupPreviewFilter].
buildRelayGroupPreviewFilter(key.groupId, since?.get(relay)?.time),
// A DM channel takes reactions/deletions the same way a group channel does.
buildRelayGroupAuxFilter(key.groupId, DmHistoryTuning.recentBoundary()),
) +
filterGroupNotificationsToPubkey(
relay = relay,
pubkey = me,
groupIds = listOf(key.groupId.id),
since = since?.get(relay)?.time,
)
}
override fun id(key: BuzzDmJoinedChatTailQueryState) = key.account
override fun id(key: BuzzDmJoinedChatTailQueryState) = key.groupId
override fun newSub(key: BuzzDmJoinedChatTailQueryState): Subscription {
windowLoad.startLoading(key.account.scope)
@@ -23,13 +23,16 @@ package com.vitorpamplona.amethyst.ui.screen.loggedIn.chats.publicChannels.relay
import androidx.compose.runtime.Composable
import androidx.compose.runtime.LaunchedEffect
import androidx.compose.runtime.getValue
import androidx.compose.runtime.key
import androidx.compose.runtime.remember
import androidx.lifecycle.compose.collectAsStateWithLifecycle
import com.vitorpamplona.amethyst.commons.model.buzz.BuzzDmChannels
import com.vitorpamplona.amethyst.commons.model.buzz.BuzzDmRegistry
import com.vitorpamplona.amethyst.commons.model.chats.ChatFeedType
import com.vitorpamplona.amethyst.commons.relayClient.subscriptions.KeyDataSourceSubscription
import com.vitorpamplona.amethyst.commons.relayClient.subscriptions.LifecycleAwareKeyDataSourceSubscription
import com.vitorpamplona.amethyst.ui.screen.loggedIn.AccountViewModel
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
import com.vitorpamplona.quartz.nip29RelayGroups.GroupId
/**
@@ -52,18 +55,35 @@ fun RelayGroupJoinedStatePreload(accountViewModel: AccountViewModel) {
/**
* Always-on preview **live tail** for joined groups' recent chat, mounted alongside [RelayGroupJoinedStatePreload]
* — keeps the Messages-list previews reflecting the true newest message app-wide. Re-derives on join/leave.
* — keeps the Messages-list previews reflecting the true newest message app-wide.
*
* Mounts **one subscription per joined group** (not one per host relay): a subscription whose `#h` carries
* more than one group id is registered as a *global* subscription by `block/buzz` and then never receives
* channel-scoped events, so a batched tail backfills once and goes deaf. See
* [buildRelayGroupJoinedChatTailFilter]. Joining/leaving adds or removes a key, so nothing needs an
* explicit invalidate; a group whose relay url won't normalize is skipped, exactly as the Messages feed
* filter skips it.
*/
@Composable
fun RelayGroupJoinedChatTailPreload(accountViewModel: AccountViewModel) {
val account = accountViewModel.account
val dataSource = accountViewModel.dataSources().relayGroupJoinedChatTail
val state = remember(account) { RelayGroupJoinedChatTailQueryState(account) }
val joined by account.relayGroupList.liveRelayGroupList.collectAsStateWithLifecycle()
LaunchedEffect(joined) { dataSource.invalidateFilters() }
val enabledFeeds by account.settings.enabledChatFeeds.collectAsStateWithLifecycle()
KeyDataSourceSubscription(state, dataSource)
if (ChatFeedType.NIP29 in enabledFeeds) {
joined.forEach { tag ->
val relay = RelayUrlNormalizer.normalizeOrNull(tag.relayUrl)
if (relay != null) {
val groupId = GroupId(tag.groupId, relay)
key(groupId) {
val state = remember(account, groupId) { RelayGroupJoinedChatTailQueryState(account, groupId) }
KeyDataSourceSubscription(state, dataSource)
}
}
}
}
}
/**
@@ -77,13 +97,27 @@ fun RelayGroupJoinedChatTailPreload(accountViewModel: AccountViewModel) {
fun BuzzDmJoinedChatTailPreload(accountViewModel: AccountViewModel) {
val account = accountViewModel.account
val dataSource = accountViewModel.dataSources().buzzDmJoinedChatTail
val state = remember(account) { BuzzDmJoinedChatTailQueryState(account) }
val me = account.userProfile().pubkeyHex
// Both values must be read (not just collected) so the snapshot system recomposes this
// preload — and re-derives the mounted set — whenever a DM is discovered or hidden.
val channels by BuzzDmChannels.flow.collectAsStateWithLifecycle()
val hidden by BuzzDmRegistry.hidden.collectAsStateWithLifecycle()
LaunchedEffect(channels, hidden) { dataSource.invalidateFilters() }
KeyDataSourceSubscription(state, dataSource)
val visible =
remember(channels, hidden, me) {
BuzzDmChannels.channelsFor(me).filterKeys { it !in BuzzDmRegistry.hiddenFor(me) }
}
// One subscription per DM channel — batching their ids into a single `#h` makes buzz treat the
// subscription as global and stop delivering live messages. See [buildRelayGroupJoinedChatTailFilter].
visible.forEach { (channelId, relay) ->
val groupId = GroupId(channelId, relay)
key(groupId) {
val state = remember(account, groupId) { BuzzDmJoinedChatTailQueryState(account, groupId) }
KeyDataSourceSubscription(state, dataSource)
}
}
}
/**
@@ -43,13 +43,16 @@ import com.vitorpamplona.quartz.nip01Core.relay.client.pool.RelayBasedFilter
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
import com.vitorpamplona.quartz.nip09Deletions.DeletionEvent
import com.vitorpamplona.quartz.nip22Comments.CommentEvent
import com.vitorpamplona.quartz.nip25Reactions.ReactionEvent
import com.vitorpamplona.quartz.nip29RelayGroups.GroupId
import com.vitorpamplona.quartz.nip29RelayGroups.metadata.GroupAdminsEvent
import com.vitorpamplona.quartz.nip29RelayGroups.metadata.GroupMembersEvent
import com.vitorpamplona.quartz.nip29RelayGroups.metadata.GroupMetadataEvent
import com.vitorpamplona.quartz.nip29RelayGroups.metadata.GroupPinnedEvent
import com.vitorpamplona.quartz.nip29RelayGroups.metadata.SupportedRolesEvent
import com.vitorpamplona.quartz.nip29RelayGroups.moderation.DeleteEventEvent
import com.vitorpamplona.quartz.nip29RelayGroups.tags.GroupIdTag
import com.vitorpamplona.quartz.nip51Lists.simpleGroupList.GroupTag
import com.vitorpamplona.quartz.nip7DThreads.ThreadEvent
@@ -147,6 +150,33 @@ val BUZZ_RELAY_GROUP_TIMELINE_EXTRA_KINDS =
/** The timeline kinds requested for every relay-group REQ (NIP-29 + Buzz; see above). */
val RELAY_GROUP_ALL_TIMELINE_KINDS = RELAY_GROUP_TIMELINE_KINDS + BUZZ_RELAY_GROUP_TIMELINE_EXTRA_KINDS
/**
* Channel **aux** kinds — overlays that modify an existing row rather than being a row: reactions (7),
* NIP-09 deletions (5) and the Buzz-native NIP-29 delete (9005).
*
* Requested `#h`-scoped on the channel subscription, which is how the Buzz reference client does it
* (`channelEventKinds` in its Flutter client bundles deletion + reaction + the message kinds into the
* one `#h` channel REQ). Amethyst otherwise only learns about reactions through the shared `#e`
* EventFinder query keyed on the ids of notes currently on screen — and that query carries no `#h`, so
* on a relay that scopes live delivery per channel it is registered as a *global* subscription and can
* never receive them live (a reaction inherits its target's channel, so it IS channel-scoped). The
* result was reaction chips that only ever appeared on a re-query, never as they happened.
*
* Kept OUT of [RELAY_GROUP_ALL_TIMELINE_KINDS] on purpose: that set feeds the backward history pager,
* whose `until` cursor walks `created_at`. Letting reactions consume a page's `limit` would advance the
* cursor past chat messages that were never delivered — the same class of cursor-skip bug that the
* unconditional-Buzz-kinds note above records.
*/
val RELAY_GROUP_AUX_KINDS =
listOf(
DeletionEvent.KIND,
ReactionEvent.KIND,
DeleteEventEvent.KIND,
)
/** How many aux events a channel replays on subscribe. Bounds the backfill; live delivery is unbounded. */
const val RELAY_GROUP_AUX_LIMIT = 500
/**
* Kinds requested on the **open channel's live tail only** — the timeline set plus the
* ephemeral kind-20002 typing indicator. Typing is scoped to the one channel on screen
@@ -208,25 +238,100 @@ fun buildRelayGroupStateFilters(
}
/**
* Recent chat of every joined group, **one `#h` filter per host relay** carrying that relay's group ids,
* bounded by a shared time floor ([sinceEpoch]) and **no per-group `limit`** — this is what lets the whole
* relay's groups batch into a single REQ and makes it reconnect-safe.
* Recent chat of **one** joined group: a single-valued `#h` filter on its host relay, bounded by a shared
* time floor ([sinceEpoch]) and **no `limit`** — a time floor bounds it, so it stays reconnect-safe (a
* reconnect re-issues one `since=window` REQ, never a page replay).
*
* ### Why one group per filter — and per *subscription*
*
* This used to batch every joined group on a relay into ONE filter carrying every group id in `#h`. That
* shape is valid NIP-01 (a multi-value tag filter is an OR), and relays serve it correctly for **stored**
* queries — which is exactly what made the bug so confusing: the boot backfill populated every group, so
* the Messages list looked right until it needed to update.
*
* But `block/buzz` (the relay behind `*.communities.buzz.xyz`) indexes each live subscription under a
* *single* channel uuid, resolved from the filters' `#h`. When two or more distinct ids appear anywhere
* across a subscription's filters, `extract_channel_id_from_filters` returns `None` and the subscription
* is registered as **global** — and global subscriptions deliberately never receive channel-scoped events
* (`fan_out_scoped`: "Global subscriptions do NOT receive channel-scoped events", guarding against leaking
* private channel content to a subscriber whose membership wasn't checked per channel). Net effect: a
* batched tail gets full history at EOSE and then goes permanently deaf, so the Messages row froze while
* the open chat — which always used a single-valued `#h` — stayed live.
*
* The resolution scans **all filters of a subscription**, so splitting into one filter per group is not
* enough: each group needs its **own subscription** (see [RelayGroupJoinedChatTailFilterAssembler], which
* keys its EOSE manager on [GroupId]). Relays advertise room for this — buzz allows `max_subscriptions:
* 1024` against `max_filters: 10` — and it sidesteps the filter cap for anyone in more than ten groups.
*/
fun buildRelayGroupJoinedChatTailFilters(
joined: Collection<GroupTag>,
fun buildRelayGroupJoinedChatTailFilter(
groupId: GroupId,
sinceEpoch: Long,
): List<RelayBasedFilter> =
byHostRelay(joined).map { (relay, ids) ->
RelayBasedFilter(
relay = relay,
filter =
Filter(
kinds = RELAY_GROUP_ALL_TIMELINE_KINDS,
tags = mapOf(GroupIdTag.TAG_NAME to ids.distinct()),
since = sinceEpoch,
),
)
}
): RelayBasedFilter =
RelayBasedFilter(
relay = groupId.relayUrl,
filter =
Filter(
kinds = RELAY_GROUP_ALL_TIMELINE_KINDS,
tags = mapOf(GroupIdTag.TAG_NAME to listOf(groupId.id)),
since = sinceEpoch,
),
)
/**
* The **newest message** of one joined channel, at any age: single-valued `#h`, `limit = 1`, and a
* `since` that is null until this relay EOSEs.
*
* This is what fills the channel's Messages-list row, and it is deliberately count-bounded rather than
* time-bounded. [buildRelayGroupJoinedChatTailFilter] floors at `now - 7 days` to keep recent chat warm
* in cache; on its own that means a channel nobody has posted in for eight days returns *zero* events,
* so `newestChatNote()` stays null and the row is stuck on its "No messages yet" placeholder forever —
* sorted to the bottom of Messages by `createdAt = 0`, indistinguishable from a channel you just joined.
*
* Every other roster-driven protocol on that screen already bounds by count for exactly this reason:
* NIP-28 asks `limit = 1` per followed channel
* ([com.vitorpamplona.amethyst.ui.screen.loggedIn.chats.rooms.datasource.filterLastMessageFollowingPublicChats]),
* Concord asks `limit = 10` per channel with no floor (`ConcordSubscriptionPlanner.channelPreviewFilters`).
* They can, because their row set comes from a list event — unlike NIP-17/NIP-04, whose *rooms* are
* discovered from the messages themselves and therefore need the backward pagers on the Messages list.
* A NIP-29 channel is roster-driven (kind 10009), so it needs one newest message per row, not a window.
*
* `since` comes from the per-relay EOSE map: null on a cold start (newest message whatever its age),
* then advancing, so a reconnect re-asks for one event rather than replaying a window.
*/
fun buildRelayGroupPreviewFilter(
groupId: GroupId,
sinceEpoch: Long?,
): RelayBasedFilter =
RelayBasedFilter(
relay = groupId.relayUrl,
filter =
Filter(
kinds = RELAY_GROUP_ALL_TIMELINE_KINDS,
tags = mapOf(GroupIdTag.TAG_NAME to listOf(groupId.id)),
since = sinceEpoch,
limit = 1,
),
)
/**
* Reactions/deletions for **every** message in one channel, `#h`-scoped on its host relay — see
* [RELAY_GROUP_AUX_KINDS]. Single-valued `#h`, so it can share a channel-scoped subscription with the
* chat tail without downgrading it.
*/
fun buildRelayGroupAuxFilter(
groupId: GroupId,
sinceEpoch: Long,
): RelayBasedFilter =
RelayBasedFilter(
relay = groupId.relayUrl,
filter =
Filter(
kinds = RELAY_GROUP_AUX_KINDS,
tags = mapOf(GroupIdTag.TAG_NAME to listOf(groupId.id)),
since = sinceEpoch,
limit = RELAY_GROUP_AUX_LIMIT,
),
)
/** The recent-chat live tail for a single open group, `#h`-scoped on its host relay. */
fun buildRelayGroupOpenChatTailFilter(
@@ -27,19 +27,21 @@ import com.vitorpamplona.amethyst.commons.relayClient.paging.WindowLoadTracker
import com.vitorpamplona.amethyst.commons.relayClient.paging.trackingListener
import com.vitorpamplona.amethyst.model.Account
import com.vitorpamplona.amethyst.service.relayClient.eoseManagers.PerUniqueIdEoseManager
import com.vitorpamplona.amethyst.service.relayClient.eoseManagers.launchChatFeedToggleObserver
import com.vitorpamplona.amethyst.service.relayClient.reqCommand.account.nip01Notifications.filterGroupNotificationsToPubkey
import com.vitorpamplona.amethyst.service.relays.SincePerRelayMap
import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient
import com.vitorpamplona.quartz.nip01Core.relay.client.pool.RelayBasedFilter
import com.vitorpamplona.quartz.nip01Core.relay.client.subscriptions.Subscription
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
import com.vitorpamplona.quartz.nip29RelayGroups.GroupId
import com.vitorpamplona.quartz.utils.TimeUtils
import kotlinx.coroutines.Job
import kotlinx.coroutines.flow.StateFlow
/** One screen's request to keep the user's joined groups' recent chat live. */
/** One joined group whose recent chat is kept live for the Messages-list preview. */
class RelayGroupJoinedChatTailQueryState(
val account: Account,
val groupId: GroupId,
)
/**
@@ -47,15 +49,17 @@ class RelayGroupJoinedChatTailQueryState(
* analog of the NIP-04 rooms-list live tail
* ([com.vitorpamplona.amethyst.ui.screen.loggedIn.chats.rooms.datasource.ChatroomListNip04SubAssembler]).
*
* One `#h`-scoped filter **per host relay** carries every joined group id on that relay at once,
* `since = recentBoundary()` and **no per-group limit** — a time floor bounds it, so unlike the fixed
* `limit=50` window it can batch all of a relay's groups into a single REQ. This keeps the Messages-list
* previews reflecting the true newest message and joined groups' recent chat warm in cache without
* **One subscription per joined group**, each a single-valued `#h` filter with `since = recentBoundary()`
* and no `limit` — a time floor bounds it, so it stays reconnect-safe. This keeps the Messages-list
* previews reflecting the true newest message, and joined groups' recent chat warm in cache, without
* opening each one. Older history is the [RelayGroupOpenChatHistorySubAssembler]'s job (below the floor).
*
* Batching by a shared time floor (rather than a per-group `limit` + shared per-relay EOSE `since`) is
* what makes this both correct — every joined group is covered the moment it joins — and reconnect-safe:
* a reconnect re-issues one `since=window` REQ per relay, never a per-group page replay.
* These used to be batched into one subscription per host relay carrying every group id in a single `#h`.
* That is valid NIP-01 and relays answer it correctly for stored queries, but it silently breaks *live*
* delivery on `block/buzz`, which can only index a subscription under one channel and downgrades a
* multi-id one to "global" — a class that by design receives no channel-scoped events. The batched tail
* therefore backfilled at EOSE and then went permanently deaf, freezing the Messages row while the open
* chat (always single-valued) stayed live. See [buildRelayGroupJoinedChatTailFilter] for the full trace.
*/
class RelayGroupJoinedChatTailFilterAssembler(
client: INostrClient,
@@ -74,7 +78,7 @@ class RelayGroupJoinedChatTailFilterAssembler(
class RelayGroupJoinedChatTailSubAssembler(
client: INostrClient,
allKeys: () -> Set<RelayGroupJoinedChatTailQueryState>,
) : PerUniqueIdEoseManager<RelayGroupJoinedChatTailQueryState, Account>(client, allKeys) {
) : PerUniqueIdEoseManager<RelayGroupJoinedChatTailQueryState, GroupId>(client, allKeys) {
private val windowLoad = WindowLoadTracker("relayGroup.preview.live")
val loadingMore: StateFlow<Boolean> = windowLoad.loading
@@ -82,42 +86,51 @@ class RelayGroupJoinedChatTailSubAssembler(
key: RelayGroupJoinedChatTailQueryState,
since: SincePerRelayMap?,
): List<RelayBasedFilter>? {
if (!key.account.settings.isChatFeedEnabled(ChatFeedType.NIP29)) {
windowLoad.setExpectedRelays(emptySet())
return null
}
val joined = key.account.relayGroupList.liveRelayGroupList.value
if (joined.isEmpty()) {
windowLoad.setExpectedRelays(emptySet())
return null
}
// The preload only mounts keys while the toggle is on; re-checked here so a flip to off
// empties any subscription that outlives the recomposition.
if (!key.account.settings.isChatFeedEnabled(ChatFeedType.NIP29)) return null
// One #h filter per host relay carrying every joined group id on it; bounded by the shared
// recent-tail floor, so no per-group limit and no per-group re-subscribe on join.
val filters = buildRelayGroupJoinedChatTailFilters(joined, DmHistoryTuning.recentBoundary())
windowLoad.setExpectedRelays(filters.mapTo(mutableSetOf()) { it.relay })
return filters
val relay = key.groupId.relayUrl
// The tracker is shared by every key of this assembler, so it must describe the whole joined
// fleet rather than whichever channel happened to rebuild last — otherwise each key would
// overwrite the expected set with its own single relay. [loadingMore] is what the Messages rows
// read to tell "still fetching this channel's newest message" apart from "channel is empty".
windowLoad.setExpectedRelays(
key.account.relayGroupList.liveRelayGroupList.value
.mapNotNullTo(mutableSetOf()) { RelayUrlNormalizer.normalizeOrNull(it.relayUrl) },
)
// Two filters, both `#h`-scoped to THIS group, so the subscription still resolves to a single
// channel on buzz (its resolver only bails when two *distinct* ids appear, or when some filter
// carries no channel tag at all). Each is queried independently, so they keep their own window:
// - the chat tail: every timeline kind, time-floored, no limit
// - group activity addressed to me: reactions/zaps/reports/replies (kinds the tail doesn't ask
// for), `#p` = me + `limit`, `since` from EOSE so a cold start is all-time and a reconnect is
// cheap. This used to ride the account-wide notifications subscription, which also carries
// inbox filters with no `#h` — that alone forces the whole subscription global on buzz, so
// live group reactions never arrived until the next launch.
return listOf(
buildRelayGroupJoinedChatTailFilter(key.groupId, DmHistoryTuning.recentBoundary()),
// The row's newest message at ANY age. The tail above floors at 7 days to keep recent chat
// warm, which on its own strands a channel quiet for longer than that on a placeholder row.
buildRelayGroupPreviewFilter(key.groupId, since?.get(relay)?.time),
// Reactions/deletions for every message in this channel — the shape Buzz's own client uses.
buildRelayGroupAuxFilter(key.groupId, DmHistoryTuning.recentBoundary()),
) +
filterGroupNotificationsToPubkey(
relay = relay,
pubkey = key.account.userProfile().pubkeyHex,
groupIds = listOf(key.groupId.id),
since = since?.get(relay)?.time,
)
}
override fun id(key: RelayGroupJoinedChatTailQueryState) = key.account
private val toggleJobs = mutableMapOf<Account, Job>()
override fun id(key: RelayGroupJoinedChatTailQueryState) = key.groupId
override fun newSub(key: RelayGroupJoinedChatTailQueryState): Subscription {
windowLoad.startLoading(key.account.scope)
toggleJobs.remove(key.account)?.cancel()
toggleJobs[key.account] =
key.account.scope.launchChatFeedToggleObserver(key.account, ChatFeedType.NIP29) { invalidateFilters() }
return requestNewSubscription(
windowLoad.trackingListener { relay: NormalizedRelayUrl, filters -> newEose(key, relay, TimeUtils.now(), filters) },
)
}
override fun endSub(
key: Account,
subId: String,
) {
super.endSub(key, subId)
toggleJobs.remove(key)?.cancel()
}
}
@@ -76,7 +76,12 @@ class RelayGroupOpenChatTailSubAssembler(
since: SincePerRelayMap?,
): List<RelayBasedFilter> {
windowLoad.setExpectedRelays(setOf(key.groupId.relayUrl))
return listOf(buildRelayGroupOpenChatTailFilter(key.groupId, DmHistoryTuning.recentBoundary()))
return listOf(
buildRelayGroupOpenChatTailFilter(key.groupId, DmHistoryTuning.recentBoundary()),
// Reactions/deletions for every message on screen, `#h`-scoped rather than `#e`-scoped —
// see [RELAY_GROUP_AUX_KINDS]. Covers a non-joined group opened by link too.
buildRelayGroupAuxFilter(key.groupId, DmHistoryTuning.recentBoundary()),
)
}
override fun id(key: RelayGroupOpenChatTailQueryState) = key.groupId
@@ -458,9 +458,19 @@ private fun RelayGroupRoomCompose(
// rather than "author: {json}". Plain chat messages fall through to the usual framing.
buzzTimelinePreviewSummary(noteEvent) ?: "$authorName: ${noteEvent.content.take(200)}"
} else {
// Event-less placeholder row for a just-joined group with no messages yet — say so
// explicitly (like Marmot groups) instead of an empty second line.
stringRes(R.string.relay_group_no_messages_yet)
// Event-less placeholder row. Until the channel's `limit = 1` preview REQ settles we cannot
// tell an empty channel from one whose newest message simply hasn't arrived, and claiming
// "No messages yet" for a busy channel reads as a bug — so say "Loading" while the joined
// fleet is still fetching, and commit to the empty wording only once it has.
val stillFetching by accountViewModel
.dataSources()
.relayGroupJoinedChatTail.tail.loadingMore
.collectAsStateWithLifecycle()
if (stillFetching) {
stringRes(R.string.loading_feed)
} else {
stringRes(R.string.relay_group_no_messages_yet)
}
}
val groupPicture = channel.profilePicture()?.ifBlank { null }
@@ -23,7 +23,6 @@ package com.vitorpamplona.amethyst.ui.screen.loggedIn.chats.publicChannels.relay
import com.vitorpamplona.quartz.buzz.stream.StreamMessageV2Event
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
import com.vitorpamplona.quartz.nip29RelayGroups.GroupId
import com.vitorpamplona.quartz.nip51Lists.simpleGroupList.GroupTag
import org.junit.Assert.assertTrue
import org.junit.Test
@@ -50,10 +49,9 @@ class BuzzTimelineKindsTest {
assertTrue(openTail.filter.kinds!!.contains(StreamMessageV2Event.KIND))
val fleet =
buildRelayGroupJoinedChatTailFilters(
listOf(GroupTag("g1", buzzRelay.url), GroupTag("g2", vanillaRelay.url)),
sinceEpoch = 0L,
)
listOf(GroupId("g1", buzzRelay), GroupId("g2", vanillaRelay)).map {
buildRelayGroupJoinedChatTailFilter(it, sinceEpoch = 0L)
}
fleet.forEach { assertTrue(it.filter.kinds!!.contains(StreamMessageV2Event.KIND)) }
}
}
@@ -24,13 +24,16 @@ import com.vitorpamplona.quartz.buzz.forum.ForumCommentEvent
import com.vitorpamplona.quartz.buzz.forum.ForumPostEvent
import com.vitorpamplona.quartz.nip01Core.relay.client.pool.FiltersChanged
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
import com.vitorpamplona.quartz.nip09Deletions.DeletionEvent
import com.vitorpamplona.quartz.nip22Comments.CommentEvent
import com.vitorpamplona.quartz.nip25Reactions.ReactionEvent
import com.vitorpamplona.quartz.nip29RelayGroups.GroupId
import com.vitorpamplona.quartz.nip29RelayGroups.metadata.GroupAdminsEvent
import com.vitorpamplona.quartz.nip29RelayGroups.metadata.GroupMembersEvent
import com.vitorpamplona.quartz.nip29RelayGroups.metadata.GroupMetadataEvent
import com.vitorpamplona.quartz.nip29RelayGroups.metadata.GroupPinnedEvent
import com.vitorpamplona.quartz.nip29RelayGroups.metadata.SupportedRolesEvent
import com.vitorpamplona.quartz.nip29RelayGroups.moderation.DeleteEventEvent
import com.vitorpamplona.quartz.nip51Lists.simpleGroupList.GroupTag
import com.vitorpamplona.quartz.nip7DThreads.ThreadEvent
import com.vitorpamplona.quartz.nip88Polls.poll.PollEvent
@@ -110,21 +113,96 @@ class RelayGroupFilterBuildersTest {
filters.filter { it.relay == relayB }.forEach { assertNull(it.filter.since) }
}
// --- Joined chat tail (always-on): batched #h per relay, time floor, NO per-group limit ---
// --- Joined chat tail (always-on): ONE single-valued #h per group, time floor, NO limit ---
@Test
fun `joined chat tail batches one h-filter per relay with no count limit`() {
val filters = buildRelayGroupJoinedChatTailFilters(joined, 999L)
assertEquals(2, filters.size)
fun `joined chat tail carries exactly one h value per group with no count limit`() {
val filter = buildRelayGroupJoinedChatTailFilter(GroupId("g1", relayA), 999L)
val a = filters.single { it.relay == relayA }
// The joined tail now always carries the Buzz timeline kinds too (see
assertEquals(relayA, filter.relay)
// The joined tail always carries the Buzz timeline kinds too (see
// RELAY_GROUP_ALL_TIMELINE_KINDS / BuzzTimelineKindsTest).
assertEquals(timelineKinds + BUZZ_RELAY_GROUP_TIMELINE_EXTRA_KINDS, a.filter.kinds)
assertEquals(setOf("g1", "g2"), a.filter.tags!!["h"]!!.toSet())
assertEquals(999L, a.filter.since)
assertNull("a time-floored tail must NOT cap by count (that is what lets it batch)", a.filter.limit)
assertNull(a.filter.until)
assertEquals(timelineKinds + BUZZ_RELAY_GROUP_TIMELINE_EXTRA_KINDS, filter.filter.kinds)
assertEquals(listOf("g1"), filter.filter.tags!!["h"])
assertEquals(999L, filter.filter.since)
assertNull("a time-floored tail must NOT cap by count", filter.filter.limit)
assertNull(filter.filter.until)
}
/**
* Regression guard for the Messages-list freeze: `block/buzz` indexes a live subscription under a
* single channel uuid resolved from `#h`, and downgrades any subscription whose filters mention two
* or more distinct ids to a *global* one — which by design receives no channel-scoped events. The
* stored query still answers correctly, so the symptom was "backfills on launch, then never updates".
* Each joined group must therefore get its own single-valued filter (and its own subscription).
*/
@Test
fun `joined chat tail never batches multiple groups into one h filter`() {
joined.forEach { tag ->
val relay = RelayUrlNormalizer.normalizeOrNull(tag.relayUrl)!!
val values = buildRelayGroupJoinedChatTailFilter(GroupId(tag.groupId, relay), 1L).filter.tags!!["h"]!!
assertEquals("a multi-value #h makes buzz register the sub as global and stop live delivery", 1, values.size)
assertEquals(tag.groupId, values.single())
}
}
/**
* The Messages-row preview is bounded by COUNT, never by time. The chat tail floors at `now - 7 days`
* to keep recent chat warm; on its own that leaves a channel quiet for longer stuck on its
* "No messages yet" placeholder forever, sorted to the bottom by `createdAt = 0`. Roster-driven
* protocols fill their rows this way — NIP-28 with `limit = 1`, Concord with `limit = 10` — because
* their row set comes from a list event and only needs one newest message per row.
*/
@Test
fun `preview filter asks for one newest message with no time floor on a cold start`() {
val filter = buildRelayGroupPreviewFilter(g1OnA, null)
assertEquals(relayA, filter.relay)
assertEquals(listOf("g1"), filter.filter.tags!!["h"])
assertEquals(1, filter.filter.limit)
assertNull("a cold start must reach back past any window to fill the row", filter.filter.since)
assertNull(filter.filter.until)
}
@Test
fun `preview filter rides the per-relay EOSE since once caught up`() {
assertEquals(555L, buildRelayGroupPreviewFilter(g1OnA, 555L).filter.since)
}
/** Single-valued `#h`, so it can share the channel-scoped subscription without downgrading it. */
@Test
fun `preview filter names exactly one channel`() {
assertEquals(1, buildRelayGroupPreviewFilter(g1OnA, null).filter.tags!!["h"]!!.size)
}
/**
* Reactions/deletions ride the channel's own `#h` subscription, the way the Buzz reference client
* does it (its `channelEventKinds` bundles deletion + reaction into the one `#h` channel REQ).
* Amethyst's other source of reactions is the shared `#e` EventFinder query, which carries no `#h`
* and is therefore registered as a *global* subscription — a class that never receives
* channel-scoped events live, and a reaction inherits its target message's channel.
*/
@Test
fun `aux filter carries reactions and deletions scoped to one group`() {
val filter = buildRelayGroupAuxFilter(GroupId("g1", relayA), 42L)
assertEquals(relayA, filter.relay)
assertEquals(listOf(DeletionEvent.KIND, ReactionEvent.KIND, DeleteEventEvent.KIND), filter.filter.kinds)
assertEquals(listOf("g1"), filter.filter.tags!!["h"])
assertEquals(42L, filter.filter.since)
assertEquals(RELAY_GROUP_AUX_LIMIT, filter.filter.limit)
}
/**
* The aux kinds must stay OUT of the timeline set: that set feeds the backward history pager, whose
* `until` cursor walks `created_at`. Reactions eating a page's `limit` would advance the cursor past
* chat messages that were never delivered.
*/
@Test
fun `aux kinds are not in the timeline set the history pager pages over`() {
RELAY_GROUP_AUX_KINDS.forEach {
assertFalse("aux kind $it must not consume history-page limit", it in RELAY_GROUP_ALL_TIMELINE_KINDS)
}
}
// --- Open chat tail: a single group's recent #h on the host relay ---
@@ -239,7 +317,6 @@ class RelayGroupFilterBuildersTest {
@Test
fun `empty joined set produces no filters`() {
assertTrue(buildRelayGroupStateFilters(emptySet()) { null }.isEmpty())
assertTrue(buildRelayGroupJoinedChatTailFilters(emptySet(), 1L).isEmpty())
}
// --- Reconnect stability: a `since`-only bump must NOT trigger a fresh REQ (no full replay) ---
@@ -265,8 +342,9 @@ class RelayGroupFilterBuildersTest {
@Test
fun `joined chat tail advancing its time floor is not a resend`() {
val earlier = buildRelayGroupJoinedChatTailFilters(joined, 1_000L).map { it.filter }
val later = buildRelayGroupJoinedChatTailFilters(joined, 2_000L).map { it.filter } // recentBoundary() crept forward
val groups = joined.map { GroupId(it.groupId, RelayUrlNormalizer.normalizeOrNull(it.relayUrl)!!) }
val earlier = groups.map { buildRelayGroupJoinedChatTailFilter(it, 1_000L).filter }
val later = groups.map { buildRelayGroupJoinedChatTailFilter(it, 2_000L).filter } // recentBoundary() crept forward
earlier.zip(later).forEach { (old, new) ->
assertFalse(
"a forward recent-tail floor bump is since-only → reconnect re-REQs the tail, not a full page",