fix(desktop): make native unread projection authoritative

Co-authored-by: Tyler Longwell <tlongwell@squareup.com>
Signed-off-by: Tyler Longwell <tlongwell@squareup.com>
This commit is contained in:
Perci
2026-08-16 04:20:55 -04:00
co-authored by Tyler Longwell
parent 8a392d87aa
commit a77c6c0679
3 changed files with 166 additions and 75 deletions
@@ -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<string> = new Set();
/** Build the native membership seed once per identity/relay scope. */
export function useObservedUnreadMembershipSeed(
scope: string,
pubkey: string | undefined,
followedRootIds: ReadonlySet<string> | undefined,
mutedChannelIds: ReadonlySet<string>,
): 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;
}
@@ -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<string>) => 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<string>();
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<string>) => {
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,
],
);
@@ -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<string>());
@@ -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<string>());
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) {