diff --git a/desktop/src/features/channels/unreadMembership.ts b/desktop/src/features/channels/unreadMembership.ts new file mode 100644 index 000000000..1983a49d1 --- /dev/null +++ b/desktop/src/features/channels/unreadMembership.ts @@ -0,0 +1,41 @@ +import * as React from "react"; +import type { ObservedUnreadMembershipSeed } from "@/shared/api/tauriObservedUnread"; +import { makeRootIdStore } from "@/features/channels/unreadRootIdStore"; + +export const participationStore = makeRootIdStore( + "buzz-thread-participation.v1", +); +export const authoredStore = makeRootIdStore("buzz-thread-authored.v1"); +export const mentionedStore = makeRootIdStore("buzz-thread-mentioned.v1"); +export const mutedStore = makeRootIdStore("buzz-thread-muted.v1"); + +const EMPTY_IDS: ReadonlySet = new Set(); + +/** Build the native membership seed once per identity/relay scope. */ +export function useObservedUnreadMembershipSeed( + scope: string, + pubkey: string | undefined, + followedRootIds: ReadonlySet | undefined, + mutedChannelIds: ReadonlySet, +): ObservedUnreadMembershipSeed { + const cached = React.useRef<{ + scope: string; + value: ObservedUnreadMembershipSeed; + } | null>(null); + if (cached.current?.scope !== scope) { + const read = (store: typeof participationStore) => + pubkey ? store.read(pubkey) : EMPTY_IDS; + cached.current = { + scope, + value: { + participatedRootIds: [...read(participationStore)], + authoredRootIds: [...read(authoredStore)], + mentionedRootIds: [...read(mentionedStore)], + followedRootIds: [...(followedRootIds ?? EMPTY_IDS)], + mutedRootIds: [...read(mutedStore)], + mutedChannelIds: [...mutedChannelIds], + }, + }; + } + return cached.current.value; +} diff --git a/desktop/src/features/channels/useObservedUnreadPersistence.ts b/desktop/src/features/channels/useObservedUnreadPersistence.ts index 6bde48c3f..6fb309fc6 100644 --- a/desktop/src/features/channels/useObservedUnreadPersistence.ts +++ b/desktop/src/features/channels/useObservedUnreadPersistence.ts @@ -1,4 +1,5 @@ import * as React from "react"; +export type { ObservedUnreadMembershipSeed } from "@/shared/api/tauriObservedUnread"; import { clearObservedUnreadStorage, deriveLatestByChannel, @@ -10,7 +11,10 @@ import { type ObservedUnreadRefs, } from "@/features/channels/observedUnreadStorage"; import { activityScopeKey } from "@/features/channels/threadActivityStorage"; -import type { ObservedUnreadEvent } from "@/features/channels/unreadChannelCounts"; +import { + recordObservedUnreadEvent, + type ObservedUnreadEvent, +} from "@/features/channels/unreadChannelCounts"; import { ingestObservedUnread, openObservedUnreadScope, @@ -32,6 +36,8 @@ export type ObservedUnreadPersistence = { ) => void; removeChannel: (channelId: string) => void; updateMembership: (kind: string, value: string, present: boolean) => void; + syncMarkers: (contextIds: Iterable) => void; + latestForChannel: (channelId: string) => number | undefined; clearAll: () => void; }; @@ -165,6 +171,7 @@ export function useObservedUnreadPersistence( [reopen], ); + // biome-ignore lint/correctness/useExhaustiveDependencies: mutable storage refs are stable containers const flushNative = React.useCallback(() => { const state = nativeRef.current; if (!state || queueRef.current.length === 0) return; @@ -201,6 +208,21 @@ export function useObservedUnreadPersistence( }).then(apply); }) .catch(() => { + // Native can fail after the caller stopped maintaining the legacy + // mirror. Seed fallback with the unacked events before scheduling its + // first write; acknowledged native history remains native-owned. + for (const { event } of queued) { + const { channelId, ...observed } = event; + recordObservedUnreadEvent( + observedUnreadEventsByChannelRef.current, + channelId, + observed, + 1_000, + ); + } + latestByChannelRef.current = deriveLatestByChannel( + observedUnreadEventsByChannelRef.current, + ); queueRef.current = [...queued, ...queueRef.current]; nativeFailedRef.current = true; nativeRef.current = null; @@ -287,31 +309,19 @@ export function useObservedUnreadPersistence( }; }, [normalizedPubkey, normalizedRelayUrl]); - // readStateVersion is the intentional invalidation signal; marker readers and - // mutable refs are sampled when that signal advances. - // biome-ignore lint/correctness/useExhaustiveDependencies: readStateVersion and readiness intentionally drive marker synchronization - React.useEffect(() => { - if (!isReadStateReady || scopeLoadedRef.current !== currentScope) return; - const state = nativeRef.current; - if (state) { - const contexts = new Set(); - for (const [ - channelId, - events, - ] of observedUnreadEventsByChannelRef.current) { - contexts.add(channelId); - for (const event of events.values()) { - contexts.add(`msg:${event.id}`); - if (event.rootId) contexts.add(`thread:${event.rootId}`); - } - } - const markers = [...contexts].map((contextId) => ({ + const syncMarkers = React.useCallback( + (contextIds: Iterable) => { + const state = nativeRef.current; + if (!state) return; + const markers = [...new Set(contextIds)].map((contextId) => ({ contextId, readAt: contextId.startsWith("thread:") || contextId.startsWith("msg:") ? getOwnTimestamp(contextId) : getEffectiveTimestamp(contextId), })); + if (markers.length === 0) return; + flushNative(); chainRef.current = chainRef.current.then(() => { const current = nativeRef.current; if ( @@ -331,16 +341,28 @@ export function useObservedUnreadPersistence( clearAll: false, }).then(apply); }); - } else if ( - pruneObservedUnreadByMarkers( - observedUnreadEventsByChannelRef.current, - latestByChannelRef.current, - getEffectiveTimestamp, - getOwnTimestamp, - ) - ) { - scheduleObservedUnreadWrite(currentScope, persistRefs.current); - optionsRef.current.onPruned?.(); + }, + [apply, flushNative, getEffectiveTimestamp, getOwnTimestamp], + ); + + // A read-state revision only needs to send the contexts that changed. Native + // observed rows remain authoritative; walking a renderer-side event mirror + // here would keep the duplicate read model that this store exists to retire. + // biome-ignore lint/correctness/useExhaustiveDependencies: readStateVersion is the intentional invalidation signal + React.useEffect(() => { + if (!isReadStateReady || scopeLoadedRef.current !== currentScope) return; + if (!nativeRef.current) { + if ( + pruneObservedUnreadByMarkers( + observedUnreadEventsByChannelRef.current, + latestByChannelRef.current, + getEffectiveTimestamp, + getOwnTimestamp, + ) + ) { + scheduleObservedUnreadWrite(currentScope, persistRefs.current); + optionsRef.current.onPruned?.(); + } } }, [readStateVersion, isReadStateReady]); @@ -439,6 +461,14 @@ export function useObservedUnreadPersistence( clearObservedUnreadStorage(normalizedPubkey ?? "", normalizedRelayUrl); } }, [currentScope, mutateClear, normalizedPubkey, normalizedRelayUrl]); + // biome-ignore lint/correctness/useExhaustiveDependencies: mutable projection/storage refs are stable containers + const latestForChannel = React.useCallback( + (channelId: string) => + nativeRef.current + ? projectionsRef.current.get(channelId)?.latest + : latestByChannelRef.current.get(channelId), + [], + ); const isScopeLoaded = React.useCallback( () => scopeLoadedRef.current === currentScope, [currentScope], @@ -457,6 +487,8 @@ export function useObservedUnreadPersistence( schedule, removeChannel, updateMembership, + syncMarkers, + latestForChannel, clearAll, }), [ @@ -466,6 +498,8 @@ export function useObservedUnreadPersistence( schedule, removeChannel, updateMembership, + syncMarkers, + latestForChannel, clearAll, ], ); diff --git a/desktop/src/features/channels/useUnreadChannels.ts b/desktop/src/features/channels/useUnreadChannels.ts index 8a8743525..40778305b 100644 --- a/desktop/src/features/channels/useUnreadChannels.ts +++ b/desktop/src/features/channels/useUnreadChannels.ts @@ -15,7 +15,6 @@ import { type ObservedUnreadEvent, } from "@/features/channels/unreadChannelCounts"; import { useReadState } from "@/features/channels/readState/useReadState"; -import { makeRootIdStore } from "@/features/channels/unreadRootIdStore"; import { forcedUnreadStore, type ForcedUnreadMap, @@ -50,6 +49,13 @@ export { writeActivityToStorage, } from "@/features/channels/threadActivityStorage"; import { useObservedUnreadPersistence } from "@/features/channels/useObservedUnreadPersistence"; +import { + authoredStore, + mentionedStore, + mutedStore, + participationStore, + useObservedUnreadMembershipSeed, +} from "@/features/channels/unreadMembership"; import { useThreadActivityPersistence } from "@/features/channels/useThreadActivityPersistence"; import { unreadCatchUp } from "@/shared/api/tauriUnreadCatchUp"; @@ -76,14 +82,6 @@ export function channelCatchUpEventKinds( : CHANNEL_MESSAGE_EVENT_KINDS; } -const participationStore = makeRootIdStore("buzz-thread-participation.v1"); -const authoredStore = makeRootIdStore("buzz-thread-authored.v1"); -// Thread roots where an external message @-mentioned the current user. The -// badge gate ORs this in so a mention recipient who never participated, -// authored, or followed still gets the thread-unread badge. -const mentionedStore = makeRootIdStore("buzz-thread-mentioned.v1"); -const mutedStore = makeRootIdStore("buzz-thread-muted.v1"); - function parseTimestamp(value: string | null | undefined) { if (!value) { return null; @@ -174,26 +172,6 @@ export function useUnreadChannels( pubkey ? forcedUnreadStore.read(pubkey) : {}, ); - // When a synced event advances a read marker (cross-device mark-as-read), - // remove from forcedUnreadRef so the dot clears immediately. - // biome-ignore lint/correctness/useExhaustiveDependencies: readStateVersion is the intentional drain trigger - React.useEffect(() => { - const advanced = drainSyncedAdvances(); - let anyNew = false; - for (const channelId of advanced) { - if (Object.hasOwn(forcedUnreadRef.current, channelId)) { - delete forcedUnreadRef.current[channelId]; - anyNew = true; - } - } - if (anyNew) { - if (pubkey) { - forcedUnreadStore.write(pubkey, forcedUnreadRef.current); - } - bumpLatestVersion(); - } - }, [readStateVersion, drainSyncedAdvances]); - // Root event IDs of threads where the current user has replied at least once. // Used to determine if thread replies should trigger unread notifications. const participatedRootIdsRef = React.useRef(new Set()); @@ -264,6 +242,13 @@ export function useUnreadChannels( bumpMembershipVersion(); }, [pubkey, relayClient, normalizedRelayUrl]); + const membershipSeed = useObservedUnreadMembershipSeed( + `${normalizedPubkey ?? ""}\u0000${normalizedRelayUrl}`, + pubkey, + options.followedRootIds, + mutedChannelIdsRef.current, + ); + // Persistence layer: hydration, pagehide flush, scope fence, write-through, marker-prune. const observedPersistence = useObservedUnreadPersistence( normalizedPubkey, @@ -276,16 +261,32 @@ export function useUnreadChannels( latestByChannelRef, { onPruned: bumpLatestVersion, - membershipSeed: { - participatedRootIds: [...participatedRootIdsRef.current], - authoredRootIds: [...authoredRootIdsRef.current], - mentionedRootIds: [...mentionedRootIdsRef.current], - followedRootIds: [...(options.followedRootIds ?? EMPTY_ROOT_IDS)], - mutedRootIds: [...mutedRootIdsRef.current], - mutedChannelIds: [...mutedChannelIdsRef.current], - }, + membershipSeed, }, ); + + // Forward only changed NIP-RS contexts; native owns the event read model. + // biome-ignore lint/correctness/useExhaustiveDependencies: readStateVersion is the intentional drain trigger + React.useEffect(() => { + const advanced = drainSyncedAdvances(); + observedPersistence.syncMarkers(advanced); + let anyNew = false; + for (const channelId of advanced) { + if ( + !channelId.startsWith("thread:") && + !channelId.startsWith("msg:") && + Object.hasOwn(forcedUnreadRef.current, channelId) + ) { + delete forcedUnreadRef.current[channelId]; + anyNew = true; + } + } + if (anyNew) { + if (pubkey) forcedUnreadStore.write(pubkey, forcedUnreadRef.current); + bumpLatestVersion(); + } + }, [readStateVersion, drainSyncedAdvances, observedPersistence, pubkey]); + const followedMembershipRef = React.useRef(new Set()); React.useEffect(() => { const desired = options.followedRootIds ?? EMPTY_ROOT_IDS; @@ -342,7 +343,7 @@ export function useUnreadChannels( } const observedLatest = topLevelOnly ? undefined - : latestByChannelRef.current.get(channelId); + : observedPersistence.latestForChannel(channelId); const { markAt, clearObserved } = resolveChannelReadMarker( readAt, observedLatest, @@ -398,6 +399,14 @@ export function useUnreadChannels( const recordUnreadEvent = React.useCallback( (channelId: string, event: ObservedUnreadEvent): boolean => { if (!observedPersistence.isScopeLoaded()) return false; + if (observedPersistence.isNative()) { + observedPersistence.schedule( + observedPersistence.currentScope, + channelId, + event, + ); + return true; + } const didRecord = recordObservedUnreadEvent( observedUnreadEventsByChannelRef.current, channelId, @@ -438,8 +447,12 @@ export function useUnreadChannels( // Fence latestByChannelRef on the scope guard — a stale live callback // during A→B drift must not write A's timestamp into B's hydrated ref. const scopeOk = observedPersistence.isScopeLoaded(); - const current = latestByChannelRef.current.get(channelId) ?? 0; - if (scopeOk && event.created_at > current) { + const current = observedPersistence.latestForChannel(channelId) ?? 0; + if ( + scopeOk && + !observedPersistence.isNative() && + event.created_at > current + ) { latestByChannelRef.current.set(channelId, event.created_at); } if (didRecordUnreadEvent || (scopeOk && event.created_at > current)) { @@ -707,8 +720,11 @@ export function useUnreadChannels( result.maxTrigger > (getEffectiveTimestamp(result.channelId) ?? 0) ) { const current = - latestByChannelRef.current.get(result.channelId) ?? 0; - if (result.maxTrigger > current) { + observedPersistence.latestForChannel(result.channelId) ?? 0; + if ( + !observedPersistence.isNative() && + result.maxTrigger > current + ) { latestByChannelRef.current.set( result.channelId, result.maxTrigger, @@ -819,8 +835,8 @@ export function useUnreadChannels( const nativeProjection = observedPersistence.isNative() ? observedPersistence.projectionsRef.current.get(channel.id) : undefined; - const unreadCount = nativeProjection - ? nativeProjection.count + const unreadCount = observedPersistence.isNative() + ? (nativeProjection?.count ?? 0) : latestByChannelRef.current.get(channel.id) === undefined ? 0 : countUnreadObservedEvents(observedEvents, readAtForObservedEvent); @@ -910,7 +926,7 @@ export function useUnreadChannels( for (const channelId of unreadChannelIdsRef.current) { delete forcedUnreadRef.current[channelId]; const unixSeconds = - latestByChannelRef.current.get(channelId) ?? + observedPersistence.latestForChannel(channelId) ?? getEffectiveTimestamp(channelId) ?? null; if (unixSeconds !== null) {