mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 03:38:23 +00:00
fix: keep Concord community unread flow alive across Refounding + dedupe emissions
Audit follow-ups on the Messages unread dot: - concordCommunityHasUnreadFlow captured the community session once, but a Refounding rebuilds a still-joined community's session in place (same id, new object). Because the row's flow is remember(communityId)-scoped it was never recreated, so the fan-in kept watching the dead session and the dot froze. Re-resolve the session on every concordSessions.revision tick instead — the same signal ConcordServerRoomCompose already uses to refresh the row's name/icon — which also covers channel-set changes from folds. - distinctUntilChanged() on both aggregated server flows so a new message that doesn't flip the has-unread boolean no longer pushes a redundant value into composition. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014ZkjdHANMKyu2EbM7eEM8C
This commit is contained in:
+24
-16
@@ -31,6 +31,7 @@ import com.vitorpamplona.quartz.nip22Comments.CommentEvent
|
||||
import kotlinx.coroutines.ExperimentalCoroutinesApi
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.combine
|
||||
import kotlinx.coroutines.flow.distinctUntilChanged
|
||||
import kotlinx.coroutines.flow.flatMapLatest
|
||||
import kotlinx.coroutines.flow.flowOf
|
||||
|
||||
@@ -72,22 +73,29 @@ fun concordChannelUnreadCountFlow(
|
||||
fun concordCommunityHasUnreadFlow(
|
||||
account: Account,
|
||||
communityId: String,
|
||||
): Flow<Boolean> {
|
||||
val session = account.concordSessions.sessionFor(communityId) ?: return flowOf(false)
|
||||
return session.state.flatMapLatest { state ->
|
||||
val channelKeys =
|
||||
state
|
||||
?.channels
|
||||
?.keys
|
||||
?.toList()
|
||||
.orEmpty()
|
||||
if (channelKeys.isEmpty()) {
|
||||
flowOf(false)
|
||||
} else {
|
||||
combine(channelKeys.map { concordChannelUnreadCountFlow(account, communityId, it) }) { counts -> counts.any { it > 0 } }
|
||||
}
|
||||
}
|
||||
}
|
||||
): Flow<Boolean> =
|
||||
// Re-resolve the session on every revision tick rather than capturing it once. A Refounding
|
||||
// rebuilds a still-joined community's session in place (same id, new object) and a fold changes
|
||||
// the channel set — both bump `revision`; capturing the session once would leave the fan-in
|
||||
// pointed at a dead session so the dot freezes. This mirrors how ConcordServerRoomCompose already
|
||||
// re-reads the row's name/icon off `revision`.
|
||||
account.concordSessions.revision
|
||||
.flatMapLatest {
|
||||
val channelKeys =
|
||||
account.concordSessions
|
||||
.sessionFor(communityId)
|
||||
?.state
|
||||
?.value
|
||||
?.channels
|
||||
?.keys
|
||||
?.toList()
|
||||
.orEmpty()
|
||||
if (channelKeys.isEmpty()) {
|
||||
flowOf(false)
|
||||
} else {
|
||||
combine(channelKeys.map { concordChannelUnreadCountFlow(account, communityId, it) }) { counts -> counts.any { it > 0 } }
|
||||
}
|
||||
}.distinctUntilChanged()
|
||||
|
||||
/**
|
||||
* True for a note the Concord channel *timeline* actually renders — the same predicate as
|
||||
|
||||
+12
-10
@@ -30,6 +30,7 @@ import com.vitorpamplona.quartz.nip29RelayGroups.isGroupChatContent
|
||||
import kotlinx.coroutines.ExperimentalCoroutinesApi
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.combine
|
||||
import kotlinx.coroutines.flow.distinctUntilChanged
|
||||
import kotlinx.coroutines.flow.flatMapLatest
|
||||
import kotlinx.coroutines.flow.flowOf
|
||||
|
||||
@@ -64,17 +65,18 @@ fun relayGroupServerHasUnreadFlow(
|
||||
account: Account,
|
||||
relay: NormalizedRelayUrl,
|
||||
): Flow<Boolean> =
|
||||
account.relayGroupList.liveRelayGroupList.flatMapLatest { tags ->
|
||||
val groupIds =
|
||||
tags.mapNotNull { tag ->
|
||||
if (RelayUrlNormalizer.normalizeOrNull(tag.relayUrl) == relay) GroupId(tag.groupId, relay) else null
|
||||
account.relayGroupList.liveRelayGroupList
|
||||
.flatMapLatest { tags ->
|
||||
val groupIds =
|
||||
tags.mapNotNull { tag ->
|
||||
if (RelayUrlNormalizer.normalizeOrNull(tag.relayUrl) == relay) GroupId(tag.groupId, relay) else null
|
||||
}
|
||||
if (groupIds.isEmpty()) {
|
||||
flowOf(false)
|
||||
} else {
|
||||
combine(groupIds.map { relayGroupChannelHasUnreadFlow(account, it) }) { perGroup -> perGroup.any { it } }
|
||||
}
|
||||
if (groupIds.isEmpty()) {
|
||||
flowOf(false)
|
||||
} else {
|
||||
combine(groupIds.map { relayGroupChannelHasUnreadFlow(account, it) }) { perGroup -> perGroup.any { it } }
|
||||
}
|
||||
}
|
||||
}.distinctUntilChanged()
|
||||
|
||||
/** Whether this group's message store holds any acceptable chat content created after [sinceSecs]. */
|
||||
private fun RelayGroupChannel.hasChatNewerThan(
|
||||
|
||||
Reference in New Issue
Block a user