From 78efe8901f3fb0ed1e76dbeaeec11355611a5137 Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 24 Sep 2026 02:17:29 +0000 Subject: [PATCH] refactor: share the per-note event observers The observeNote* family (note, event, reactions, zaps, reposts, replies, OTS, references, modifications, community approval) moves to commonsUI relayClient/event. Each takes the UserFinderAccount and the EventFinderFilterAssembler, defaulting to LocalUserFinderAccount and LocalEventFinder, so a shared composable needs only the note. The app's EventObservers keep their signatures and delegate with the account and data source from AccountViewModel, so call sites are unchanged and Activity roots without the locals keep working. isMinichatReply (pure quartz) moves to commons model/chats, and its test to commonTest with kotlin.test. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_0168wY9t7i9NC5u3svyMxEz6 --- .../vitorpamplona/amethyst/model/Account.kt | 2 +- .../reqCommand/event/EventObservers.kt | 401 ++------------ .../chats/minichat/MinichatFeedViewModel.kt | 2 +- .../publicChannels/dal/ChannelFeedFilter.kt | 2 +- .../relayGroup/RelayGroupUnread.kt | 2 +- .../loggedIn/redirect/LoadRedirectScreen.kt | 2 +- .../commons/model}/chats/MinichatReply.kt | 2 +- .../commons/model}/chats/MinichatReplyTest.kt | 8 +- .../relayClient/event/EventObservers.kt | 516 ++++++++++++++++++ 9 files changed, 570 insertions(+), 367 deletions(-) rename {amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn => commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model}/chats/MinichatReply.kt (98%) rename {amethyst/src/test/java/com/vitorpamplona/amethyst/ui/screen/loggedIn => commons/src/commonTest/kotlin/com/vitorpamplona/amethyst/commons/model}/chats/MinichatReplyTest.kt (96%) create mode 100644 commonsUI/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/relayClient/event/EventObservers.kt diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/model/Account.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/model/Account.kt index 3ba243658a..85b79bd3c9 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/model/Account.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/model/Account.kt @@ -2048,7 +2048,7 @@ class Account( // blessed status. 40002 survives only as a read-compat tail from the // 10002 -> 40001 -> 40002 migration, so we were the last active writer of a kind // their clients no longer thread on. Reading 40002 stays supported (see - // [com.vitorpamplona.amethyst.ui.screen.loggedIn.chats.isMinichatReply]). + // [com.vitorpamplona.amethyst.commons.model.chats.isMinichatReply]). // // Attached media rides as URLs appended to the content. val root = rootEvent.tags.buzzThreadRoot() ?: rootEvent.tags.buzzThreadReply() ?: rootEvent.id diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/relayClient/reqCommand/event/EventObservers.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/relayClient/reqCommand/event/EventObservers.kt index cde6322ad4..33e77422c6 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/relayClient/reqCommand/event/EventObservers.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/relayClient/reqCommand/event/EventObservers.kt @@ -22,465 +22,152 @@ package com.vitorpamplona.amethyst.service.relayClient.reqCommand.event import androidx.compose.runtime.Composable import androidx.compose.runtime.State -import androidx.compose.runtime.produceState -import androidx.compose.runtime.remember -import androidx.lifecycle.compose.collectAsStateWithLifecycle import com.vitorpamplona.amethyst.commons.model.Note import com.vitorpamplona.amethyst.commons.model.NoteState import com.vitorpamplona.amethyst.commons.model.User -import com.vitorpamplona.amethyst.commons.model.textNoteModifications import com.vitorpamplona.amethyst.ui.screen.loggedIn.AccountViewModel -import com.vitorpamplona.amethyst.ui.screen.loggedIn.chats.isMinichatReply import com.vitorpamplona.quartz.nip01Core.core.Event -import com.vitorpamplona.quartz.nip18Reposts.GenericRepostEvent -import com.vitorpamplona.quartz.nip18Reposts.RepostEvent -import com.vitorpamplona.quartz.nip72ModCommunities.approval.CommunityPostApprovalEvent -import com.vitorpamplona.quartz.nip72ModCommunities.definition.CommunityDefinitionEvent -import com.vitorpamplona.quartz.nip72ModCommunities.isForCommunity -import kotlinx.coroutines.Dispatchers -import kotlinx.coroutines.ExperimentalCoroutinesApi -import kotlinx.coroutines.FlowPreview -import kotlinx.coroutines.flow.combine -import kotlinx.coroutines.flow.distinctUntilChanged -import kotlinx.coroutines.flow.flowOn -import kotlinx.coroutines.flow.mapLatest -import kotlinx.coroutines.flow.sample +import com.vitorpamplona.amethyst.commons.relayClient.event.observeCommunityApprovalNeedStatus as sharedObserveCommunityApprovalNeedStatus +import com.vitorpamplona.amethyst.commons.relayClient.event.observeNote as sharedObserveNote +import com.vitorpamplona.amethyst.commons.relayClient.event.observeNoteAndMap as sharedObserveNoteAndMap +import com.vitorpamplona.amethyst.commons.relayClient.event.observeNoteEvent as sharedObserveNoteEvent +import com.vitorpamplona.amethyst.commons.relayClient.event.observeNoteEventAndMap as sharedObserveNoteEventAndMap +import com.vitorpamplona.amethyst.commons.relayClient.event.observeNoteEventAndMapNotNull as sharedObserveNoteEventAndMapNotNull +import com.vitorpamplona.amethyst.commons.relayClient.event.observeNoteHasEvent as sharedObserveNoteHasEvent +import com.vitorpamplona.amethyst.commons.relayClient.event.observeNoteMinichatReplyCount as sharedObserveNoteMinichatReplyCount +import com.vitorpamplona.amethyst.commons.relayClient.event.observeNoteModifications as sharedObserveNoteModifications +import com.vitorpamplona.amethyst.commons.relayClient.event.observeNoteOts as sharedObserveNoteOts +import com.vitorpamplona.amethyst.commons.relayClient.event.observeNoteReactionCount as sharedObserveNoteReactionCount +import com.vitorpamplona.amethyst.commons.relayClient.event.observeNoteReactions as sharedObserveNoteReactions +import com.vitorpamplona.amethyst.commons.relayClient.event.observeNoteReferences as sharedObserveNoteReferences +import com.vitorpamplona.amethyst.commons.relayClient.event.observeNoteReplies as sharedObserveNoteReplies +import com.vitorpamplona.amethyst.commons.relayClient.event.observeNoteReplyCount as sharedObserveNoteReplyCount +import com.vitorpamplona.amethyst.commons.relayClient.event.observeNoteRepostCount as sharedObserveNoteRepostCount +import com.vitorpamplona.amethyst.commons.relayClient.event.observeNoteReposts as sharedObserveNoteReposts +import com.vitorpamplona.amethyst.commons.relayClient.event.observeNoteRepostsBy as sharedObserveNoteRepostsBy +import com.vitorpamplona.amethyst.commons.relayClient.event.observeNoteZaps as sharedObserveNoteZaps + +/* + * Android overloads of the shared per-note observers in commons `relayClient/event`: they take the + * account and event-finder data source from [AccountViewModel] instead of the composition locals, + * so they also work under Activity roots that don't provide those locals. + */ @Composable fun observeNote( note: Note, accountViewModel: AccountViewModel, -): State { - // Subscribe in the relay for changes in this note. - EventFinderFilterAssemblerSubscription(note, accountViewModel) +): State = sharedObserveNote(note, accountViewModel.account, accountViewModel.dataSources().eventFinder) - // Subscribe in the LocalCache for changes that arrive in the device - val flow = remember(note) { note.flow().metadata.stateFlow } - return flow.collectAsStateWithLifecycle() -} - -/** - * [observeNote] without the relay half: watches LocalCache and asks no relay for the note. - * - * For a NIP-17 rumor, which has no fetchable id — putting one in a REQ would tell relays the - * private event's identity, the leak [com.vitorpamplona.amethyst.commons.model.Note.isPrivateRumor] - * guards everywhere else. Nothing is lost by not asking: a rumor only ever reaches the cache by - * unwrapping the envelope that carried it, and the always-on gift-wrap tail re-fetches a week of - * those on every cold start, so this flow fires on its own once the envelope lands. - */ -@Composable -fun observeNoteLocally(note: Note): State { - val flow = remember(note) { note.flow().metadata.stateFlow } - return flow.collectAsStateWithLifecycle() -} - -@OptIn(ExperimentalCoroutinesApi::class) @Composable inline fun observeNoteEvent( note: Note, accountViewModel: AccountViewModel, -): State { - // Subscribe in the relay for changes in this note. - EventFinderFilterAssemblerSubscription(note, accountViewModel) +): State = sharedObserveNoteEvent(note, accountViewModel.account, accountViewModel.dataSources().eventFinder) - // Subscribe in the LocalCache for changes that arrive in the device - val flow = - remember(note) { - note - .flow() - .metadata.stateFlow - .mapLatest { it.note.event as? T? } - } - - return flow.collectAsStateWithLifecycle(note.event as? T?) -} - -@OptIn(ExperimentalCoroutinesApi::class) @Composable fun observeNoteAndMap( note: Note, accountViewModel: AccountViewModel, map: (Note) -> T, -): State { - // Subscribe in the relay for changes in this note. - EventFinderFilterAssemblerSubscription(note, accountViewModel) +): State = sharedObserveNoteAndMap(note, accountViewModel.account, accountViewModel.dataSources().eventFinder, map) - val flow = - remember(note) { - note - .flow() - .metadata.stateFlow - .mapLatest { map(it.note) } - .distinctUntilChanged() - .flowOn(Dispatchers.IO) - } - - // Subscribe in the LocalCache for changes that arrive in the device - return flow.collectAsStateWithLifecycle(map(note)) -} - -@OptIn(ExperimentalCoroutinesApi::class) -@Suppress("UNCHECKED_CAST") @Composable fun observeNoteEventAndMapNotNull( note: Note, accountViewModel: AccountViewModel, map: (T) -> U, -): State { - // Subscribe in the relay for changes in this note. - EventFinderFilterAssemblerSubscription(note, accountViewModel) +): State = sharedObserveNoteEventAndMapNotNull(note, accountViewModel.account, accountViewModel.dataSources().eventFinder, map) - // Subscribe in the LocalCache for changes that arrive in the device - val flow = - remember(note) { - note - .flow() - .metadata.stateFlow - .mapLatest { (it.note.event as? T)?.let { map(it) } } - .distinctUntilChanged() - .flowOn(Dispatchers.IO) - } - - // Subscribe in the LocalCache for changes that arrive in the device - return flow.collectAsStateWithLifecycle((note.event as? T)?.let { map(it) }) -} - -@OptIn(ExperimentalCoroutinesApi::class) -@Suppress("UNCHECKED_CAST") @Composable fun observeNoteEventAndMap( note: Note, accountViewModel: AccountViewModel, map: (T?) -> U, -): State { - // Subscribe in the relay for changes in this note. - EventFinderFilterAssemblerSubscription(note, accountViewModel) +): State = sharedObserveNoteEventAndMap(note, accountViewModel.account, accountViewModel.dataSources().eventFinder, map) - // Subscribe in the LocalCache for changes that arrive in the device - val flow = - remember(note) { - note - .flow() - .metadata.stateFlow - .mapLatest { map(it.note.event as? T) } - .distinctUntilChanged() - .flowOn(Dispatchers.IO) - } - - // Subscribe in the LocalCache for changes that arrive in the device - return flow.collectAsStateWithLifecycle(map(note.event as? T)) -} - -@OptIn(ExperimentalCoroutinesApi::class) @Composable fun observeNoteHasEvent( note: Note, accountViewModel: AccountViewModel, -): State { - // Subscribe in the relay for changes in this note. - EventFinderFilterAssemblerSubscription(note, accountViewModel) - - // Subscribe in the LocalCache for changes that arrive in the device - val flow = - remember(note) { - note - .flow() - .metadata.stateFlow - .mapLatest { it.note.event != null } - .distinctUntilChanged() - } - - return flow.collectAsStateWithLifecycle(note.event != null) -} +): State = sharedObserveNoteHasEvent(note, accountViewModel.account, accountViewModel.dataSources().eventFinder) @Composable fun observeNoteReplies( note: Note, accountViewModel: AccountViewModel, -): State { - // Subscribe in the relay for changes in this note. - EventFinderFilterAssemblerSubscription(note, accountViewModel) +): State = sharedObserveNoteReplies(note, accountViewModel.account, accountViewModel.dataSources().eventFinder) - // Subscribe in the LocalCache for changes that arrive in the device - val flow = remember(note) { note.flow().replies.stateFlow } - return flow.collectAsStateWithLifecycle() -} - -@OptIn(ExperimentalCoroutinesApi::class, FlowPreview::class) @Composable fun observeNoteReplyCount( note: Note, accountViewModel: AccountViewModel, -): State { - // Subscribe in the relay for changes in this note. - EventFinderFilterAssemblerSubscription(note, accountViewModel) +): State = sharedObserveNoteReplyCount(note, accountViewModel.account, accountViewModel.dataSources().eventFinder) - // Subscribe in the LocalCache for changes that arrive in the device - val flow = - remember(note) { - note - .flow() - .replies.stateFlow - .sample(200) - .mapLatest { it.note.replies.size } - .distinctUntilChanged() - } - - return flow.collectAsStateWithLifecycle(note.replies.size) -} - -/** - * Count of a chat message's **minichat** replies — its kind-1111 [CommentEvent] - * children only (inline quote-replies are ordinary kind-9/42 messages and are not - * counted here). Drives the "N replies" chip that opens the minichat. - * - * Mounting this registers the message with [EventFinderFilterAssemblerSubscription], which - * batches the visible messages' ids into shared REQs for their replies (kind-1111 among - * them) — so for public chats (NIP-28/NIP-29) the thread replies load, and the chip appears, - * just by rendering the rows. Concord's kind-1111 replies instead arrive over the channel - * plane, so that REQ finds nothing there and is a harmless no-op. - */ -@OptIn(ExperimentalCoroutinesApi::class, FlowPreview::class) @Composable fun observeNoteMinichatReplyCount( note: Note, accountViewModel: AccountViewModel, -): State { - EventFinderFilterAssemblerSubscription(note, accountViewModel) - - val flow = - remember(note) { - note - .flow() - .replies.stateFlow - .sample(200) - .mapLatest { it.note.replies.count { reply -> isMinichatReply(reply.event) } } - .distinctUntilChanged() - } - - return flow.collectAsStateWithLifecycle(note.replies.count { isMinichatReply(it.event) }) -} +): State = sharedObserveNoteMinichatReplyCount(note, accountViewModel.account, accountViewModel.dataSources().eventFinder) @Composable fun observeNoteReactions( note: Note, accountViewModel: AccountViewModel, -): State { - // Subscribe in the relay for changes in this note. - EventFinderFilterAssemblerSubscription(note, accountViewModel) +): State = sharedObserveNoteReactions(note, accountViewModel.account, accountViewModel.dataSources().eventFinder) - // Subscribe in the LocalCache for changes that arrive in the device - val flow = remember(note) { note.flow().reactions.stateFlow } - return flow.collectAsStateWithLifecycle() -} - -@OptIn(ExperimentalCoroutinesApi::class, FlowPreview::class) @Composable fun observeNoteReactionCount( note: Note, accountViewModel: AccountViewModel, -): State { - // Subscribe in the relay for changes in this note. - EventFinderFilterAssemblerSubscription(note, accountViewModel) - - // Subscribe in the LocalCache for changes that arrive in the device - val flow = - remember(note) { - note - .flow() - .reactions.stateFlow - .sample(200) - .mapLatest { it.note.countReactions() } - .distinctUntilChanged() - .flowOn(Dispatchers.IO) - } - - // Subscribe in the LocalCache for changes that arrive in the device - return flow.collectAsStateWithLifecycle(note.countReactions()) -} +): State = sharedObserveNoteReactionCount(note, accountViewModel.account, accountViewModel.dataSources().eventFinder) @Composable fun observeNoteZaps( note: Note, accountViewModel: AccountViewModel, -): State { - // Subscribe in the relay for changes in this note. - EventFinderFilterAssemblerSubscription(note, accountViewModel) - - // Subscribe in the LocalCache for changes that arrive in the device - val flow = remember(note) { note.flow().zaps.stateFlow } - return flow.collectAsStateWithLifecycle() -} +): State = sharedObserveNoteZaps(note, accountViewModel.account, accountViewModel.dataSources().eventFinder) @Composable fun observeNoteReposts( note: Note, accountViewModel: AccountViewModel, -): State { - // Subscribe in the relay for changes in this note. - EventFinderFilterAssemblerSubscription(note, accountViewModel) +): State = sharedObserveNoteReposts(note, accountViewModel.account, accountViewModel.dataSources().eventFinder) - // Subscribe in the LocalCache for changes that arrive in the device - val flow = remember(note) { note.flow().boosts.stateFlow } - return flow.collectAsStateWithLifecycle() -} - -@OptIn(ExperimentalCoroutinesApi::class) @Composable fun observeNoteRepostsBy( note: Note, user: User, accountViewModel: AccountViewModel, -): State { - // Subscribe in the relay for changes in this note. - EventFinderFilterAssemblerSubscription(note, accountViewModel) +): State = sharedObserveNoteRepostsBy(note, user, accountViewModel.account, accountViewModel.dataSources().eventFinder) - // Subscribe in the LocalCache for changes that arrive in the device - val flow = - remember(note) { - note - .flow() - .boosts.stateFlow - .mapLatest { it.note.isBoostedBy(user) } - .distinctUntilChanged() - .flowOn(Dispatchers.IO) - } - - return flow.collectAsStateWithLifecycle(note.isBoostedBy(user)) -} - -@OptIn(ExperimentalCoroutinesApi::class, FlowPreview::class) @Composable fun observeNoteRepostCount( note: Note, accountViewModel: AccountViewModel, -): State { - // Subscribe in the relay for changes in this note. - EventFinderFilterAssemblerSubscription(note, accountViewModel) - - // Subscribe in the LocalCache for changes that arrive in the device - val flow = - remember(note) { - note - .flow() - .boosts.stateFlow - .sample(200) - .mapLatest { note.boosts.size } - .distinctUntilChanged() - } - - return flow.collectAsStateWithLifecycle(note.boosts.size) -} +): State = sharedObserveNoteRepostCount(note, accountViewModel.account, accountViewModel.dataSources().eventFinder) @Composable fun observeNoteReferences( note: Note, accountViewModel: AccountViewModel, -): State { - // Subscribe in the relay for changes in this note. - EventFinderFilterAssemblerSubscription(note, accountViewModel) - - // Subscribe in the LocalCache for changes that arrive in the device. - // On-chain zaps piggyback `note.flow().zaps.stateFlow` because Note.addOnchainZap - // invalidates the same `flowSet.zaps`. If that ever moves to a dedicated flow, - // add it here too — otherwise on-chain-only notes will stop pinging the chevron. - val flow = - remember(note) { - combine( - note.flow().zaps.stateFlow, - note.flow().boosts.stateFlow, - note.flow().reactions.stateFlow, - ) { zapState, _, _ -> - zapState.note.hasZapsBoostsOrReactions() - }.distinctUntilChanged() - } - - return flow.collectAsStateWithLifecycle(note.hasZapsBoostsOrReactions()) -} +): State = sharedObserveNoteReferences(note, accountViewModel.account, accountViewModel.dataSources().eventFinder) @Composable fun observeNoteOts( note: Note, accountViewModel: AccountViewModel, -): State { - // Subscribe in the relay for changes in this note. - EventFinderFilterAssemblerSubscription(note, accountViewModel) +): State = sharedObserveNoteOts(note, accountViewModel.account, accountViewModel.dataSources().eventFinder) - // Subscribe in the LocalCache for changes that arrive in the device - val flow = remember(note) { note.flow().ots.stateFlow } - return flow.collectAsStateWithLifecycle() -} - -// Resolves the actual modification list off the main thread and filters identical results, -// so the caller's LaunchedEffect only fires when the list of edits truly changes. -// `sample(500)` collapses bursts — a heavily-edited note can emit hundreds of times during -// initial relay sync, and we only need the last state per ~half second. -// Returns `null` until the first IO resolution completes — callers should treat that as -// "still loading" and not flip their UI to "no edits". -@OptIn(ExperimentalCoroutinesApi::class, FlowPreview::class) @Composable fun observeNoteModifications( note: Note, accountViewModel: AccountViewModel, -): State?> { - // Subscribe in the relay for changes in this note. - EventFinderFilterAssemblerSubscription(note, accountViewModel) - - return produceState?>(initialValue = null, note) { - note - .flow() - .edits - .stateFlow - .sample(500) - .mapLatest { note.textNoteModifications() } - .distinctUntilChanged() - .flowOn(Dispatchers.IO) - .collect { value = it } - } -} +): State?> = sharedObserveNoteModifications(note, accountViewModel.account, accountViewModel.dataSources().eventFinder) @Composable fun observeCommunityApprovalNeedStatus( note: Note, community: Note, accountViewModel: AccountViewModel, -): State { - // Subscribe in the relay for changes in this note. - EventFinderFilterAssemblerSubscription(note, accountViewModel) - - // Subscribe in the LocalCache for changes that arrive in the device - val flow = - remember(note, community) { - combine( - community.flow().metadata.stateFlow, - note.flow().boosts.stateFlow, - ) { communityMetadata, boosts -> - (communityMetadata.note.event as? CommunityDefinitionEvent)?.let { communityDefEvent -> - val moderators = communityDefEvent.moderatorKeys().toSet() - - if (note.author?.pubkeyHex in moderators) { - false - } else { - val isModerator = accountViewModel.account.userProfile().pubkeyHex in moderators - - if (isModerator) { - val wasAlreadyApproved = - note.boosts.any { - val approvalEvent = it.event - (approvalEvent is CommunityPostApprovalEvent || approvalEvent is RepostEvent || approvalEvent is GenericRepostEvent) && - approvalEvent.pubKey in moderators && - approvalEvent.isForCommunity(community.idHex) - } - !wasAlreadyApproved - } else { - false - } - } - } - }.distinctUntilChanged() - .flowOn(Dispatchers.IO) - } - - // Subscribe in the LocalCache for changes that arrive in the device - return flow.collectAsStateWithLifecycle(false) -} +): State = sharedObserveCommunityApprovalNeedStatus(note, community, accountViewModel.account, accountViewModel.dataSources().eventFinder) diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/minichat/MinichatFeedViewModel.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/minichat/MinichatFeedViewModel.kt index 4c261853bd..4d26fdc2fd 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/minichat/MinichatFeedViewModel.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/minichat/MinichatFeedViewModel.kt @@ -24,8 +24,8 @@ import androidx.lifecycle.ViewModel import androidx.lifecycle.ViewModelProvider import androidx.lifecycle.viewModelScope import com.vitorpamplona.amethyst.commons.model.Note +import com.vitorpamplona.amethyst.commons.model.chats.isMinichatReply import com.vitorpamplona.amethyst.model.Account -import com.vitorpamplona.amethyst.ui.screen.loggedIn.chats.isMinichatReply import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.ExperimentalCoroutinesApi import kotlinx.coroutines.flow.SharingStarted diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/publicChannels/dal/ChannelFeedFilter.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/publicChannels/dal/ChannelFeedFilter.kt index b2f6b23c37..f656839204 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/publicChannels/dal/ChannelFeedFilter.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/publicChannels/dal/ChannelFeedFilter.kt @@ -24,9 +24,9 @@ import com.vitorpamplona.amethyst.commons.feeds.AdditiveFeedFilter import com.vitorpamplona.amethyst.commons.feeds.ChangesFlowFilter import com.vitorpamplona.amethyst.commons.model.Channel import com.vitorpamplona.amethyst.commons.model.Note +import com.vitorpamplona.amethyst.commons.model.chats.isMinichatReply import com.vitorpamplona.amethyst.model.Account import com.vitorpamplona.amethyst.ui.dal.sortedByDefaultFeedOrder -import com.vitorpamplona.amethyst.ui.screen.loggedIn.chats.isMinichatReply class ChannelFeedFilter( val channel: Channel, diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/publicChannels/relayGroup/RelayGroupUnread.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/publicChannels/relayGroup/RelayGroupUnread.kt index 37105e8a5c..3937effabe 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/publicChannels/relayGroup/RelayGroupUnread.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/publicChannels/relayGroup/RelayGroupUnread.kt @@ -22,10 +22,10 @@ package com.vitorpamplona.amethyst.ui.screen.loggedIn.chats.publicChannels.relay import com.vitorpamplona.amethyst.commons.model.Note import com.vitorpamplona.amethyst.commons.model.cache.LocalCache +import com.vitorpamplona.amethyst.commons.model.chats.isMinichatReply import com.vitorpamplona.amethyst.commons.model.nip29RelayGroups.RelayGroupChannel import com.vitorpamplona.amethyst.model.Account import com.vitorpamplona.amethyst.ui.dal.sortedByDefaultFeedOrder -import com.vitorpamplona.amethyst.ui.screen.loggedIn.chats.isMinichatReply import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer import com.vitorpamplona.quartz.nip29RelayGroups.GroupId diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/redirect/LoadRedirectScreen.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/redirect/LoadRedirectScreen.kt index 5c25c80a54..a2060f9e78 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/redirect/LoadRedirectScreen.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/redirect/LoadRedirectScreen.kt @@ -36,6 +36,7 @@ import androidx.compose.ui.text.style.TextAlign import androidx.compose.ui.unit.dp import com.vitorpamplona.amethyst.commons.model.Note import com.vitorpamplona.amethyst.commons.model.navigation.Route +import com.vitorpamplona.amethyst.commons.relayClient.event.observeNoteLocally import com.vitorpamplona.amethyst.commons.resources.Res import com.vitorpamplona.amethyst.commons.resources.looking_for_event import com.vitorpamplona.amethyst.commons.resources.looking_for_event_title @@ -44,7 +45,6 @@ import com.vitorpamplona.amethyst.commons.ui.navigation.navs.INav import com.vitorpamplona.amethyst.commons.ui.navigation.topbars.TopBarExtensibleWithBackButton import com.vitorpamplona.amethyst.commons.ui.stringRes import com.vitorpamplona.amethyst.service.relayClient.reqCommand.event.observeNote -import com.vitorpamplona.amethyst.service.relayClient.reqCommand.event.observeNoteLocally import com.vitorpamplona.amethyst.ui.components.LoadNote import com.vitorpamplona.amethyst.ui.layouts.DisappearingScaffold import com.vitorpamplona.amethyst.ui.navigation.navs.Nav diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/MinichatReply.kt b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/chats/MinichatReply.kt similarity index 98% rename from amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/MinichatReply.kt rename to commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/chats/MinichatReply.kt index c6b8d8d4f9..1e97f7e0a5 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/MinichatReply.kt +++ b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/chats/MinichatReply.kt @@ -18,7 +18,7 @@ * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. */ -package com.vitorpamplona.amethyst.ui.screen.loggedIn.chats +package com.vitorpamplona.amethyst.commons.model.chats import com.vitorpamplona.quartz.buzz.stream.StreamMessageV2Event import com.vitorpamplona.quartz.buzz.threading.buzzThreadReply diff --git a/amethyst/src/test/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/MinichatReplyTest.kt b/commons/src/commonTest/kotlin/com/vitorpamplona/amethyst/commons/model/chats/MinichatReplyTest.kt similarity index 96% rename from amethyst/src/test/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/MinichatReplyTest.kt rename to commons/src/commonTest/kotlin/com/vitorpamplona/amethyst/commons/model/chats/MinichatReplyTest.kt index acb9de2a98..27e0a33c10 100644 --- a/amethyst/src/test/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/MinichatReplyTest.kt +++ b/commons/src/commonTest/kotlin/com/vitorpamplona/amethyst/commons/model/chats/MinichatReplyTest.kt @@ -18,12 +18,12 @@ * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. */ -package com.vitorpamplona.amethyst.ui.screen.loggedIn.chats +package com.vitorpamplona.amethyst.commons.model.chats import com.vitorpamplona.quartz.nipC7Chats.ChatEvent -import org.junit.Assert.assertFalse -import org.junit.Assert.assertTrue -import org.junit.Test +import kotlin.test.Test +import kotlin.test.assertFalse +import kotlin.test.assertTrue /** * A kind-9 chat message is a thread reply only when its `e` tag carries a NIP-10 marker. diff --git a/commonsUI/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/relayClient/event/EventObservers.kt b/commonsUI/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/relayClient/event/EventObservers.kt new file mode 100644 index 0000000000..d105d02d99 --- /dev/null +++ b/commonsUI/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/relayClient/event/EventObservers.kt @@ -0,0 +1,516 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.amethyst.commons.relayClient.event + +import androidx.compose.runtime.Composable +import androidx.compose.runtime.State +import androidx.compose.runtime.produceState +import androidx.compose.runtime.remember +import androidx.lifecycle.compose.collectAsStateWithLifecycle +import com.vitorpamplona.amethyst.commons.model.Note +import com.vitorpamplona.amethyst.commons.model.NoteState +import com.vitorpamplona.amethyst.commons.model.User +import com.vitorpamplona.amethyst.commons.model.chats.isMinichatReply +import com.vitorpamplona.amethyst.commons.model.textNoteModifications +import com.vitorpamplona.amethyst.commons.relayClient.user.LocalUserFinderAccount +import com.vitorpamplona.amethyst.commons.relayClient.user.UserFinderAccount +import com.vitorpamplona.quartz.nip01Core.core.Event +import com.vitorpamplona.quartz.nip18Reposts.GenericRepostEvent +import com.vitorpamplona.quartz.nip18Reposts.RepostEvent +import com.vitorpamplona.quartz.nip72ModCommunities.approval.CommunityPostApprovalEvent +import com.vitorpamplona.quartz.nip72ModCommunities.definition.CommunityDefinitionEvent +import com.vitorpamplona.quartz.nip72ModCommunities.isForCommunity +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.ExperimentalCoroutinesApi +import kotlinx.coroutines.FlowPreview +import kotlinx.coroutines.IO +import kotlinx.coroutines.flow.combine +import kotlinx.coroutines.flow.distinctUntilChanged +import kotlinx.coroutines.flow.flowOn +import kotlinx.coroutines.flow.mapLatest +import kotlinx.coroutines.flow.sample + +/* + * Per-note observers: each subscribes [note] with the relay event finder (batched with every other + * on-screen note) and exposes a LocalCache flow of the note as Compose state. + * + * [account] and [dataSource] default to the front end's [LocalUserFinderAccount] and + * [LocalEventFinder], so a shared composable needs only the note. Callers outside a composition + * that provides them (a separate Activity root) pass both explicitly. + */ + +@Composable +fun observeNote( + note: Note, + account: UserFinderAccount = LocalUserFinderAccount.current, + dataSource: EventFinderFilterAssembler = LocalEventFinder.current, +): State { + // Subscribe in the relay for changes in this note. + EventFinderFilterAssemblerSubscription(note, account, dataSource) + + // Subscribe in the LocalCache for changes that arrive in the device + val flow = remember(note) { note.flow().metadata.stateFlow } + return flow.collectAsStateWithLifecycle() +} + +/** + * [observeNote] without the relay half: watches LocalCache and asks no relay for the note. + * + * For a NIP-17 rumor, which has no fetchable id — putting one in a REQ would tell relays the + * private event's identity, the leak [com.vitorpamplona.amethyst.commons.model.Note.isPrivateRumor] + * guards everywhere else. Nothing is lost by not asking: a rumor only ever reaches the cache by + * unwrapping the envelope that carried it, and the always-on gift-wrap tail re-fetches a week of + * those on every cold start, so this flow fires on its own once the envelope lands. + */ +@Composable +fun observeNoteLocally(note: Note): State { + val flow = remember(note) { note.flow().metadata.stateFlow } + return flow.collectAsStateWithLifecycle() +} + +@OptIn(ExperimentalCoroutinesApi::class) +@Composable +inline fun observeNoteEvent( + note: Note, + account: UserFinderAccount = LocalUserFinderAccount.current, + dataSource: EventFinderFilterAssembler = LocalEventFinder.current, +): State { + // Subscribe in the relay for changes in this note. + EventFinderFilterAssemblerSubscription(note, account, dataSource) + + // Subscribe in the LocalCache for changes that arrive in the device + val flow = + remember(note) { + note + .flow() + .metadata.stateFlow + .mapLatest { it.note.event as? T? } + } + + return flow.collectAsStateWithLifecycle(note.event as? T?) +} + +@OptIn(ExperimentalCoroutinesApi::class) +@Composable +fun observeNoteAndMap( + note: Note, + account: UserFinderAccount = LocalUserFinderAccount.current, + dataSource: EventFinderFilterAssembler = LocalEventFinder.current, + map: (Note) -> T, +): State { + // Subscribe in the relay for changes in this note. + EventFinderFilterAssemblerSubscription(note, account, dataSource) + + val flow = + remember(note) { + note + .flow() + .metadata.stateFlow + .mapLatest { map(it.note) } + .distinctUntilChanged() + .flowOn(Dispatchers.IO) + } + + // Subscribe in the LocalCache for changes that arrive in the device + return flow.collectAsStateWithLifecycle(map(note)) +} + +@OptIn(ExperimentalCoroutinesApi::class) +@Suppress("UNCHECKED_CAST") +@Composable +fun observeNoteEventAndMapNotNull( + note: Note, + account: UserFinderAccount = LocalUserFinderAccount.current, + dataSource: EventFinderFilterAssembler = LocalEventFinder.current, + map: (T) -> U, +): State { + // Subscribe in the relay for changes in this note. + EventFinderFilterAssemblerSubscription(note, account, dataSource) + + // Subscribe in the LocalCache for changes that arrive in the device + val flow = + remember(note) { + note + .flow() + .metadata.stateFlow + .mapLatest { (it.note.event as? T)?.let { map(it) } } + .distinctUntilChanged() + .flowOn(Dispatchers.IO) + } + + // Subscribe in the LocalCache for changes that arrive in the device + return flow.collectAsStateWithLifecycle((note.event as? T)?.let { map(it) }) +} + +@OptIn(ExperimentalCoroutinesApi::class) +@Suppress("UNCHECKED_CAST") +@Composable +fun observeNoteEventAndMap( + note: Note, + account: UserFinderAccount = LocalUserFinderAccount.current, + dataSource: EventFinderFilterAssembler = LocalEventFinder.current, + map: (T?) -> U, +): State { + // Subscribe in the relay for changes in this note. + EventFinderFilterAssemblerSubscription(note, account, dataSource) + + // Subscribe in the LocalCache for changes that arrive in the device + val flow = + remember(note) { + note + .flow() + .metadata.stateFlow + .mapLatest { map(it.note.event as? T) } + .distinctUntilChanged() + .flowOn(Dispatchers.IO) + } + + // Subscribe in the LocalCache for changes that arrive in the device + return flow.collectAsStateWithLifecycle(map(note.event as? T)) +} + +@OptIn(ExperimentalCoroutinesApi::class) +@Composable +fun observeNoteHasEvent( + note: Note, + account: UserFinderAccount = LocalUserFinderAccount.current, + dataSource: EventFinderFilterAssembler = LocalEventFinder.current, +): State { + // Subscribe in the relay for changes in this note. + EventFinderFilterAssemblerSubscription(note, account, dataSource) + + // Subscribe in the LocalCache for changes that arrive in the device + val flow = + remember(note) { + note + .flow() + .metadata.stateFlow + .mapLatest { it.note.event != null } + .distinctUntilChanged() + } + + return flow.collectAsStateWithLifecycle(note.event != null) +} + +@Composable +fun observeNoteReplies( + note: Note, + account: UserFinderAccount = LocalUserFinderAccount.current, + dataSource: EventFinderFilterAssembler = LocalEventFinder.current, +): State { + // Subscribe in the relay for changes in this note. + EventFinderFilterAssemblerSubscription(note, account, dataSource) + + // Subscribe in the LocalCache for changes that arrive in the device + val flow = remember(note) { note.flow().replies.stateFlow } + return flow.collectAsStateWithLifecycle() +} + +@OptIn(ExperimentalCoroutinesApi::class, FlowPreview::class) +@Composable +fun observeNoteReplyCount( + note: Note, + account: UserFinderAccount = LocalUserFinderAccount.current, + dataSource: EventFinderFilterAssembler = LocalEventFinder.current, +): State { + // Subscribe in the relay for changes in this note. + EventFinderFilterAssemblerSubscription(note, account, dataSource) + + // Subscribe in the LocalCache for changes that arrive in the device + val flow = + remember(note) { + note + .flow() + .replies.stateFlow + .sample(200) + .mapLatest { it.note.replies.size } + .distinctUntilChanged() + } + + return flow.collectAsStateWithLifecycle(note.replies.size) +} + +/** + * Count of a chat message's **minichat** replies — its kind-1111 [CommentEvent] + * children only (inline quote-replies are ordinary kind-9/42 messages and are not + * counted here). Drives the "N replies" chip that opens the minichat. + * + * Mounting this registers the message with [EventFinderFilterAssemblerSubscription], which + * batches the visible messages' ids into shared REQs for their replies (kind-1111 among + * them) — so for public chats (NIP-28/NIP-29) the thread replies load, and the chip appears, + * just by rendering the rows. Concord's kind-1111 replies instead arrive over the channel + * plane, so that REQ finds nothing there and is a harmless no-op. + */ +@OptIn(ExperimentalCoroutinesApi::class, FlowPreview::class) +@Composable +fun observeNoteMinichatReplyCount( + note: Note, + account: UserFinderAccount = LocalUserFinderAccount.current, + dataSource: EventFinderFilterAssembler = LocalEventFinder.current, +): State { + EventFinderFilterAssemblerSubscription(note, account, dataSource) + + val flow = + remember(note) { + note + .flow() + .replies.stateFlow + .sample(200) + .mapLatest { it.note.replies.count { reply -> isMinichatReply(reply.event) } } + .distinctUntilChanged() + } + + return flow.collectAsStateWithLifecycle(note.replies.count { isMinichatReply(it.event) }) +} + +@Composable +fun observeNoteReactions( + note: Note, + account: UserFinderAccount = LocalUserFinderAccount.current, + dataSource: EventFinderFilterAssembler = LocalEventFinder.current, +): State { + // Subscribe in the relay for changes in this note. + EventFinderFilterAssemblerSubscription(note, account, dataSource) + + // Subscribe in the LocalCache for changes that arrive in the device + val flow = remember(note) { note.flow().reactions.stateFlow } + return flow.collectAsStateWithLifecycle() +} + +@OptIn(ExperimentalCoroutinesApi::class, FlowPreview::class) +@Composable +fun observeNoteReactionCount( + note: Note, + account: UserFinderAccount = LocalUserFinderAccount.current, + dataSource: EventFinderFilterAssembler = LocalEventFinder.current, +): State { + // Subscribe in the relay for changes in this note. + EventFinderFilterAssemblerSubscription(note, account, dataSource) + + // Subscribe in the LocalCache for changes that arrive in the device + val flow = + remember(note) { + note + .flow() + .reactions.stateFlow + .sample(200) + .mapLatest { it.note.countReactions() } + .distinctUntilChanged() + .flowOn(Dispatchers.IO) + } + + // Subscribe in the LocalCache for changes that arrive in the device + return flow.collectAsStateWithLifecycle(note.countReactions()) +} + +@Composable +fun observeNoteZaps( + note: Note, + account: UserFinderAccount = LocalUserFinderAccount.current, + dataSource: EventFinderFilterAssembler = LocalEventFinder.current, +): State { + // Subscribe in the relay for changes in this note. + EventFinderFilterAssemblerSubscription(note, account, dataSource) + + // Subscribe in the LocalCache for changes that arrive in the device + val flow = remember(note) { note.flow().zaps.stateFlow } + return flow.collectAsStateWithLifecycle() +} + +@Composable +fun observeNoteReposts( + note: Note, + account: UserFinderAccount = LocalUserFinderAccount.current, + dataSource: EventFinderFilterAssembler = LocalEventFinder.current, +): State { + // Subscribe in the relay for changes in this note. + EventFinderFilterAssemblerSubscription(note, account, dataSource) + + // Subscribe in the LocalCache for changes that arrive in the device + val flow = remember(note) { note.flow().boosts.stateFlow } + return flow.collectAsStateWithLifecycle() +} + +@OptIn(ExperimentalCoroutinesApi::class) +@Composable +fun observeNoteRepostsBy( + note: Note, + user: User, + account: UserFinderAccount = LocalUserFinderAccount.current, + dataSource: EventFinderFilterAssembler = LocalEventFinder.current, +): State { + // Subscribe in the relay for changes in this note. + EventFinderFilterAssemblerSubscription(note, account, dataSource) + + // Subscribe in the LocalCache for changes that arrive in the device + val flow = + remember(note) { + note + .flow() + .boosts.stateFlow + .mapLatest { it.note.isBoostedBy(user) } + .distinctUntilChanged() + .flowOn(Dispatchers.IO) + } + + return flow.collectAsStateWithLifecycle(note.isBoostedBy(user)) +} + +@OptIn(ExperimentalCoroutinesApi::class, FlowPreview::class) +@Composable +fun observeNoteRepostCount( + note: Note, + account: UserFinderAccount = LocalUserFinderAccount.current, + dataSource: EventFinderFilterAssembler = LocalEventFinder.current, +): State { + // Subscribe in the relay for changes in this note. + EventFinderFilterAssemblerSubscription(note, account, dataSource) + + // Subscribe in the LocalCache for changes that arrive in the device + val flow = + remember(note) { + note + .flow() + .boosts.stateFlow + .sample(200) + .mapLatest { note.boosts.size } + .distinctUntilChanged() + } + + return flow.collectAsStateWithLifecycle(note.boosts.size) +} + +@Composable +fun observeNoteReferences( + note: Note, + account: UserFinderAccount = LocalUserFinderAccount.current, + dataSource: EventFinderFilterAssembler = LocalEventFinder.current, +): State { + // Subscribe in the relay for changes in this note. + EventFinderFilterAssemblerSubscription(note, account, dataSource) + + // Subscribe in the LocalCache for changes that arrive in the device. + // On-chain zaps piggyback `note.flow().zaps.stateFlow` because Note.addOnchainZap + // invalidates the same `flowSet.zaps`. If that ever moves to a dedicated flow, + // add it here too — otherwise on-chain-only notes will stop pinging the chevron. + val flow = + remember(note) { + combine( + note.flow().zaps.stateFlow, + note.flow().boosts.stateFlow, + note.flow().reactions.stateFlow, + ) { zapState, _, _ -> + zapState.note.hasZapsBoostsOrReactions() + }.distinctUntilChanged() + } + + return flow.collectAsStateWithLifecycle(note.hasZapsBoostsOrReactions()) +} + +@Composable +fun observeNoteOts( + note: Note, + account: UserFinderAccount = LocalUserFinderAccount.current, + dataSource: EventFinderFilterAssembler = LocalEventFinder.current, +): State { + // Subscribe in the relay for changes in this note. + EventFinderFilterAssemblerSubscription(note, account, dataSource) + + // Subscribe in the LocalCache for changes that arrive in the device + val flow = remember(note) { note.flow().ots.stateFlow } + return flow.collectAsStateWithLifecycle() +} + +// Resolves the actual modification list off the main thread and filters identical results, +// so the caller's LaunchedEffect only fires when the list of edits truly changes. +// `sample(500)` collapses bursts — a heavily-edited note can emit hundreds of times during +// initial relay sync, and we only need the last state per ~half second. +// Returns `null` until the first IO resolution completes — callers should treat that as +// "still loading" and not flip their UI to "no edits". +@OptIn(ExperimentalCoroutinesApi::class, FlowPreview::class) +@Composable +fun observeNoteModifications( + note: Note, + account: UserFinderAccount = LocalUserFinderAccount.current, + dataSource: EventFinderFilterAssembler = LocalEventFinder.current, +): State?> { + // Subscribe in the relay for changes in this note. + EventFinderFilterAssemblerSubscription(note, account, dataSource) + + return produceState?>(initialValue = null, note) { + note + .flow() + .edits + .stateFlow + .sample(500) + .mapLatest { note.textNoteModifications() } + .distinctUntilChanged() + .flowOn(Dispatchers.IO) + .collect { value = it } + } +} + +@Composable +fun observeCommunityApprovalNeedStatus( + note: Note, + community: Note, + account: UserFinderAccount = LocalUserFinderAccount.current, + dataSource: EventFinderFilterAssembler = LocalEventFinder.current, +): State { + // Subscribe in the relay for changes in this note. + EventFinderFilterAssemblerSubscription(note, account, dataSource) + + // Subscribe in the LocalCache for changes that arrive in the device + val flow = + remember(note, community) { + combine( + community.flow().metadata.stateFlow, + note.flow().boosts.stateFlow, + ) { communityMetadata, boosts -> + (communityMetadata.note.event as? CommunityDefinitionEvent)?.let { communityDefEvent -> + val moderators = communityDefEvent.moderatorKeys().toSet() + + if (note.author?.pubkeyHex in moderators) { + false + } else { + val isModerator = account.userFinderPubkeyHex in moderators + + if (isModerator) { + val wasAlreadyApproved = + note.boosts.any { + val approvalEvent = it.event + (approvalEvent is CommunityPostApprovalEvent || approvalEvent is RepostEvent || approvalEvent is GenericRepostEvent) && + approvalEvent.pubKey in moderators && + approvalEvent.isForCommunity(community.idHex) + } + !wasAlreadyApproved + } else { + false + } + } + } + }.distinctUntilChanged() + .flowOn(Dispatchers.IO) + } + + // Subscribe in the LocalCache for changes that arrive in the device + return flow.collectAsStateWithLifecycle(false) +}