fix: source the heartbeat outbox fetcher from the ungated cache scan

The outbox fetcher read its announcement set from the gated DVM feed list —
a death spiral: a DVM dropped for a stale beat left the fetch batch, its
beats were never fetched again, and the drop became permanent. The fetcher
could only ever help DVMs that were already visible.

The assembler now takes cache-backed sources from the front end:
LocalCache.cachedDvmAnnouncements (every cached k=5300 announcement,
newest-first, capped at 100 — the gate's other eligibility checks, WITHOUT
the freshness gate), a per-author relay lookup unioning the NIP-65 outbox
with the cached relay hints (same mix as the event finder), and 31990/NIP-65
observation flows as re-issue drivers.
This commit is contained in:
Dr. Tobias Baur
2026-09-10 15:09:57 +02:00
parent 6a02d787b6
commit cd67c56909
6 changed files with 138 additions and 44 deletions
@@ -102,12 +102,19 @@ the session-scoped watcher for pinned DVMs. Traffic is negligible (a few pinned
reach the *user's* discovery relays — but DVMs publish beats to their own write relays, and
relays don't gossip, so alive DVMs whose beats never overlap the user's relay set stayed
invisible (their detail screens proved the beats existed on the outbox). `DiscoveryDvmHeartbeatSubAssembler`
joins the discovery assembler group and, while Discover is composed, batches the DVM list's
announcement authors per **DVM outbox relay** (`kinds = [11998], authors = [those pubkeys],
since = now - 420`, coverage-ranked and capped at 12 relays; authors with unknown outboxes
rely on the global REQ as fallback). It re-issues when the DVM list's membership changes and
unwraps `FeedState.Loaded` to the inner feed flow (the wrapper is reused, so only the inner
flow emits real list changes).
joins the discovery assembler group and, while Discover is composed, batches the cached
content-discovery announcements' authors per **DVM outbox relay** (`kinds = [11998],
authors = [those pubkeys], since = now - 420`, coverage-ranked and capped at 12 relays;
authors with unknown outboxes/hints rely on the global REQ as fallback). It re-issues when
the cached announcement set or the NIP-65 relay lists move.
The announcement source MUST be the **ungated cache scan**
(`LocalCache.cachedDvmAnnouncements` — every cached k=5300 announcement, newest first, capped
at 100), not the gated feed list. Sourcing from the gated list is a death spiral: a DVM
leaves the gated list the moment its beat ages out, the fetcher would stop covering it, and
no beat would ever arrive to bring it back — any transient staleness becomes a permanent
drop. The relay lookup unions the author's NIP-65 outbox with the cached relay hints for the
author (the same mix the event finder's `potentialRelaysToFindAddress` uses).
## 6. Invalidation — closing the two silent gaps
@@ -20,8 +20,10 @@
*/
package com.vitorpamplona.amethyst.model
import com.vitorpamplona.amethyst.commons.model.cache.filterIntoSet
import com.vitorpamplona.quartz.nip01Core.core.Address
import com.vitorpamplona.quartz.nip89AppHandlers.definition.AppDefinitionEvent
import com.vitorpamplona.quartz.nip90Dvms.contentDiscoveryRequest.NIP90ContentDiscoveryRequestEvent
import com.vitorpamplona.quartz.nip90Dvms.dvmHeartbeat.DvmHeartbeatEvent
import com.vitorpamplona.quartz.utils.TimeUtils
@@ -33,3 +35,20 @@ fun LocalCache.hasFreshDvmHeartbeat(
appDef: AppDefinitionEvent,
now: Long = TimeUtils.now(),
): Boolean = dvmHeartbeatOf(appDef)?.isFreshAt(now) == true
/**
* Every cached content-discovery announcement, WITHOUT the freshness gate — this is the source the
* heartbeat outbox fetcher must use. Sourcing from the gated feed list would drop a DVM the moment
* its beat went stale, remove it from the fetch batch, and make the drop permanent (the fetcher
* could only ever help DVMs that were already visible). Applies the gate's other eligibility
* checks (a real content-discovery DVM, not a paid subscription app), newest first, capped.
*/
fun LocalCache.cachedDvmAnnouncements(limit: Int = 100): List<AppDefinitionEvent> =
addressables
.filterIntoSet(AppDefinitionEvent.KIND) { _, note ->
(note.event as? AppDefinitionEvent)?.let {
it.appMetaData()?.subscription != true && it.includeKind(NIP90ContentDiscoveryRequestEvent.KIND)
} == true
}.mapNotNull { it.event as? AppDefinitionEvent }
.sortedByDescending { it.createdAt }
.take(limit)
@@ -31,6 +31,7 @@ import com.vitorpamplona.amethyst.commons.relayClient.chess.ChessFilterAssembler
import com.vitorpamplona.amethyst.commons.relayClient.communities.CommunityFilterAssembler
import com.vitorpamplona.amethyst.commons.relayClient.communities.list.CommunitiesListFilterAssembler
import com.vitorpamplona.amethyst.commons.relayClient.discover.DiscoveryFilterAssembler
import com.vitorpamplona.amethyst.commons.relayClient.discover.nip90DVMs.DvmHeartbeatSources
import com.vitorpamplona.amethyst.commons.relayClient.emojipacks.BrowseEmojiSetsFilterAssembler
import com.vitorpamplona.amethyst.commons.relayClient.event.EventFinderFilterAssembler
import com.vitorpamplona.amethyst.commons.relayClient.followPacks.FollowPacksFilterAssembler
@@ -69,6 +70,7 @@ import com.vitorpamplona.amethyst.commons.relayClient.video.VideoFilterAssembler
import com.vitorpamplona.amethyst.commons.relayClient.wallet.OnchainZapsFilterAssembler
import com.vitorpamplona.amethyst.commons.relayClient.workouts.WorkoutsFilterAssembler
import com.vitorpamplona.amethyst.model.LocalCache
import com.vitorpamplona.amethyst.model.cachedDvmAnnouncements
import com.vitorpamplona.amethyst.service.relayClient.reqCommand.account.AccountFilterAssembler
import com.vitorpamplona.amethyst.service.relayClient.reqCommand.account.AccountForegroundFilterAssembler
import com.vitorpamplona.amethyst.service.relayClient.reqCommand.channel.ChannelFinderFilterAssemblyGroup
@@ -94,6 +96,9 @@ import com.vitorpamplona.amethyst.ui.screen.loggedIn.nests.datasource.NestRoomLi
import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.RelayOfflineTracker
import com.vitorpamplona.quartz.nip01Core.relay.client.auth.IAuthStatus
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.nip65RelayList.AdvertisedRelayListEvent
import com.vitorpamplona.quartz.nip89AppHandlers.definition.AppDefinitionEvent
import kotlinx.coroutines.CoroutineScope
class RelaySubscriptionsCoordinator(
@@ -113,7 +118,25 @@ class RelaySubscriptionsCoordinator(
val home = HomeFilterAssembler(client)
val chatroomList = ChatroomListFilterAssembler(client)
val video = VideoFilterAssembler(client)
val discovery = DiscoveryFilterAssembler(client, cache)
val discovery =
DiscoveryFilterAssembler(
client,
dvmHeartbeat =
DvmHeartbeatSources(
announcements = cache::cachedDvmAnnouncements,
outboxRelaysFor = { pubkey ->
buildSet {
cache.getUserIfExists(pubkey)?.outboxRelays()?.let { addAll(it) }
addAll(cache.relayHints.hintsForKey(pubkey))
}
},
changes =
listOf(
cache.observeNotes(Filter(kinds = listOf(AppDefinitionEvent.KIND))),
cache.observeNotes(Filter(kinds = listOf(AdvertisedRelayListEvent.KIND))),
),
),
)
// loaders of content that is not yet in the device.
// they are active when looking at events, users, channels.
@@ -120,4 +120,51 @@ class DvmHeartbeatTest {
fun noHeartbeatMeansNoLiveness() {
assertFalse(LocalCache.hasFreshDvmHeartbeat(appDef("dvm-never"), TimeUtils.now()))
}
@Test
fun theUngatedAnnouncementScanKeepsDvmsTheGateWouldHide() {
// The outbox fetcher must source announcements from the cache, NOT from the gated feed
// list: a DVM dropped for a stale beat must keep receiving outbox beats or it can never
// come back. Subscription apps and non-content-discovery apps stay excluded.
val alive = appDef("scan-dvm")
val subscriptionApp =
AppDefinitionEvent(
id = "c0".repeat(32),
pubKey = "ab".repeat(32),
createdAt = 1_760_000_500L,
tags = arrayOf(arrayOf("d", "subs"), arrayOf("k", "5300")),
content = """{"name":"Paid","subscription":true}""",
sig = "cc".repeat(64),
)
val nonDiscoveryApp =
AppDefinitionEvent(
id = "c1".repeat(32),
pubKey = "cb".repeat(32),
createdAt = 1_760_000_100L,
tags = arrayOf(arrayOf("d", "other"), arrayOf("k", "9999")),
content = """{"name":"Other"}""",
sig = "cc".repeat(64),
)
LocalCache.justConsume(appDef("dvm-x"), null, true)
LocalCache.justConsume(nonDiscoveryApp, null, true)
val consumed =
AppDefinitionEvent(
id = "c2".repeat(32),
pubKey = appDefPubKey,
createdAt = 1_760_000_000L,
tags = arrayOf(arrayOf("d", "subs2"), arrayOf("k", "5300")),
content = """{"name":"Paid2","subscription":true}""",
sig = "cc".repeat(64),
)
LocalCache.justConsume(consumed, null, true)
val scanned = LocalCache.cachedDvmAnnouncements()
assertTrue("dvm-x is not yet visible (no beat) but must still be sourced", scanned.any { it.dTag() == "dvm-x" })
assertFalse("subscription apps are not content-discovery DVMs", scanned.any { it.dTag() == "subs" })
assertFalse("k=9999 apps are not content-discovery DVMs", scanned.any { it.dTag() == "other" })
assertEquals("newest-first so the cap keeps the most relevant announcements", scanned.sortedByDescending { it.createdAt }, scanned)
assertTrue("capped", scanned.size <= 100)
}
}
@@ -21,10 +21,10 @@
package com.vitorpamplona.amethyst.commons.relayClient.discover
import com.vitorpamplona.amethyst.commons.model.IAccount
import com.vitorpamplona.amethyst.commons.model.cache.ICacheProvider
import com.vitorpamplona.amethyst.commons.model.topNavFeeds.IFeedTopNavPerRelayFilterSet
import com.vitorpamplona.amethyst.commons.model.topNavFeeds.TopFilter
import com.vitorpamplona.amethyst.commons.relayClient.discover.nip90DVMs.DiscoveryDvmHeartbeatSubAssembler
import com.vitorpamplona.amethyst.commons.relayClient.discover.nip90DVMs.DvmHeartbeatSources
import com.vitorpamplona.amethyst.commons.relayClient.topNavFeeds.TopNavFeedFilterAssembler
import com.vitorpamplona.amethyst.commons.relayClient.topNavFeeds.TopNavFeedQueryState
import com.vitorpamplona.amethyst.commons.ui.feeds.FeedContentState
@@ -58,12 +58,12 @@ class DiscoveryQueryState(
class DiscoveryFilterAssembler(
client: INostrClient,
cache: ICacheProvider,
dvmHeartbeat: DvmHeartbeatSources,
) : TopNavFeedFilterAssembler<DiscoveryQueryState>({ keys ->
listOf(
DiscoveryLongFormClassifiedsAndDVMSubAssembler1(client, keys),
DiscoveryFollowsSetsAndLiveStreamsSubAssembler2(client, keys),
DiscoveryPublicChatsAndCommunitiesSubAssembler3(client, keys),
DiscoveryDvmHeartbeatSubAssembler(client, cache, keys),
DiscoveryDvmHeartbeatSubAssembler(client, keys, dvmHeartbeat),
)
})
@@ -20,13 +20,11 @@
*/
package com.vitorpamplona.amethyst.commons.relayClient.discover.nip90DVMs
import com.vitorpamplona.amethyst.commons.model.cache.ICacheProvider
import com.vitorpamplona.amethyst.commons.relayClient.discover.DiscoveryQueryState
import com.vitorpamplona.amethyst.commons.relayClient.subscriptions.ExplainedFilter
import com.vitorpamplona.amethyst.commons.relayClient.subscriptions.SubPurpose
import com.vitorpamplona.amethyst.commons.relayClient.topNavFeeds.TopNavFeedSubAssembler
import com.vitorpamplona.amethyst.commons.relays.SincePerRelayMap
import com.vitorpamplona.amethyst.commons.ui.feeds.FeedState
import com.vitorpamplona.quartz.nip01Core.core.HexKey
import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient
import com.vitorpamplona.quartz.nip01Core.relay.client.pool.RelayBasedFilter
@@ -34,11 +32,8 @@ import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.nip89AppHandlers.definition.AppDefinitionEvent
import com.vitorpamplona.quartz.nip90Dvms.dvmHeartbeat.DvmHeartbeatEvent
import com.vitorpamplona.quartz.utils.TimeUtils
import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.flatMapLatest
import kotlinx.coroutines.flow.flowOf
/** How many DVM outbox relays the fetcher may open at once; coverage-ranked, so the top relays carry most authors. */
private const val MAX_OUTBOX_RELAYS = 12
@@ -94,47 +89,50 @@ fun dvmHeartbeatOutboxFilters(
}
/**
* The discovery-side subscription that runs [dvmHeartbeatOutboxFilters] for the DVM list's current
* announcement set, alive while the Discover screen is composed (it joins the same assembler group
* and lifecycle as the other discovery sub-assemblers).
* Cache-backed inputs the outbox fetcher needs, provided by the front end (the cache query and
* relay-hint surface are platform caches, not commons).
*
* Re-issues when the DVM list's membership changes: [FeedState.Loaded] reuses its wrapper, so the
* invalidator unwraps to the inner feed flow, which emits on every real list change. No floor
* collectors ([floors] is empty) — the announcement set is the driver, not note timestamps.
* [announcements] MUST be the ungated announcement set (every cached content-discovery DVM). The
* gated feed list would turn any transient staleness into a permanent drop: a DVM leaves the
* gated list the moment its beat ages out, the fetcher would stop covering it, and no beat would
* ever arrive to bring it back.
*
* [outboxRelaysFor] resolves where a DVM publishes its beats — NIP-65 outbox relays plus any
* relay hints for the author (the same mix the event finder uses).
*
* [changes] drive re-issues: when the cached announcement set or the outbox data moves, the
* batches are recomputed.
*/
class DvmHeartbeatSources(
val announcements: () -> List<AppDefinitionEvent>,
val outboxRelaysFor: (HexKey) -> Collection<NormalizedRelayUrl>,
val changes: List<Flow<*>>,
)
/**
* The discovery-side subscription that runs [dvmHeartbeatOutboxFilters] for the cached
* announcement set, alive while the Discover screen is composed (it joins the same assembler
* group and lifecycle as the other discovery sub-assemblers).
*
* No floor collectors ([floors] is empty) — the announcement set and outbox data ([changes]) are
* the drivers, not note timestamps.
*/
class DiscoveryDvmHeartbeatSubAssembler(
client: INostrClient,
private val cache: ICacheProvider,
allKeys: () -> Set<DiscoveryQueryState>,
private val sources: DvmHeartbeatSources,
) : TopNavFeedSubAssembler<DiscoveryQueryState>(client, allKeys) {
override fun updateFilter(
key: DiscoveryQueryState,
since: SincePerRelayMap?,
): List<RelayBasedFilter> {
val announcements =
(key.dvms.feedContent.value as? FeedState.Loaded)
?.feed
?.value
?.list
?.mapNotNull { it.event as? AppDefinitionEvent }
.orEmpty()
return dvmHeartbeatOutboxFilters(
announcements,
outboxRelaysFor = { pubkey ->
cache.getUserIfExists(pubkey)?.outboxRelays().orEmpty()
},
): List<RelayBasedFilter> =
dvmHeartbeatOutboxFilters(
sources.announcements(),
outboxRelaysFor = sources.outboxRelaysFor,
now = TimeUtils.now(),
)
}
override fun floors(key: DiscoveryQueryState): List<StateFlow<Long?>> = emptyList()
@OptIn(ExperimentalCoroutinesApi::class)
override fun extraInvalidators(key: DiscoveryQueryState): List<Flow<*>> =
listOf(
key.dvms.feedContent.flatMapLatest { state ->
(state as? FeedState.Loaded)?.feed ?: flowOf(null)
},
)
override fun extraInvalidators(key: DiscoveryQueryState): List<Flow<*>> = sources.changes
}