fix convo list cache updates from chat log events

- use continue instead of return when a convo isn't found in cache, so
  the rest of the event batch isn't dropped (the bus advances its
  cursor past the batch, so dropped logs are never redelivered)
- on logAcceptConvo, flip status to accepted in every cache holding the
  convo, including 'all'-status caches like the always-mounted unread
  query; previously the stale copy kept status 'request' and the next
  logCreateMessage could resurrect the convo in the requests inbox
- add a rev guard so log-driven updates skip when the log isn't newer
  than the cached convo, preventing stale refetch snapshots and log
  events from double-applying against each other

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
Samuel Newman
2026-06-12 10:36:53 +03:00
parent e475a1ca36
commit 40dcf8d853
+239 -136
View File
@@ -282,8 +282,11 @@ export function ListConvosProviderInner({
mutateMembers(convoId, list =>
list.some(m => m.did === did) ? list : list.concat(newMember),
)
mutateConvoView(convoId, convo =>
addMemberToConvoView(convo, newMember, rev, alreadyKnownMember),
mutateConvoView(
convoId,
withRevGuard(rev, convo =>
addMemberToConvoView(convo, newMember, rev, alreadyKnownMember),
),
)
}
@@ -301,8 +304,11 @@ export function ListConvosProviderInner({
>(listConvoMembersQueryKey(convoId))
?.some(m => m.did === did) === false
mutateMembers(convoId, list => list.filter(m => m.did !== did))
mutateConvoView(convoId, convo =>
removeMemberFromConvoView(convo, did, rev, alreadyRemovedMember),
mutateConvoView(
convoId,
withRevGuard(rev, convo =>
removeMemberFromConvoView(convo, did, rev, alreadyRemovedMember),
),
)
}
@@ -317,24 +323,27 @@ export function ListConvosProviderInner({
// link preview so its viewer state reflects the lost membership.
void invalidateJoinLinkPreviewsForConvo(queryClient, log.convoId)
} else if (ChatBskyConvoDefs.isLogDeleteMessage(log)) {
updateConvoInAllLists(log.convoId, convo => {
if (
(ChatBskyConvoDefs.isDeletedMessageView(log.message) ||
ChatBskyConvoDefs.isMessageView(log.message)) &&
(ChatBskyConvoDefs.isDeletedMessageView(convo.lastMessage) ||
ChatBskyConvoDefs.isMessageView(convo.lastMessage))
) {
return log.message.id === convo.lastMessage.id
? {
...convo,
rev: log.rev,
lastMessage: log.message,
}
: convo
} else {
return convo
}
})
updateConvoInAllLists(
log.convoId,
withRevGuard(log.rev, convo => {
if (
(ChatBskyConvoDefs.isDeletedMessageView(log.message) ||
ChatBskyConvoDefs.isMessageView(log.message)) &&
(ChatBskyConvoDefs.isDeletedMessageView(convo.lastMessage) ||
ChatBskyConvoDefs.isMessageView(convo.lastMessage))
) {
return log.message.id === convo.lastMessage.id
? {
...convo,
rev: log.rev,
lastMessage: log.message,
}
: convo
} else {
return convo
}
}),
)
} else if (ChatBskyConvoDefs.isLogCreateMessage(log)) {
// Store in a new var to avoid TS errors due to closures.
const logRef: ChatBskyConvoDefs.LogCreateMessage = log
@@ -356,9 +365,20 @@ export function ListConvosProviderInner({
}
if (!foundConvo) {
// Convo not found, trigger refetch
// Convo not found, trigger refetch. Use continue (not return) so
// the remaining logs in this batch still apply - the bus advances
// its cursor past this batch, so a dropped log is never
// redelivered.
debouncedRefetch()
return
continue
}
// Rev guard. updatedConvo is built once from foundConvo and applied
// across caches, so guarding here (rather than per-cache) is both
// simplest and correct - skip if the log isn't newer than the
// cached convo.
if (logRef.rev <= foundConvo.rev) {
continue
}
// add relatedProfiles to members list, but making sure to dedupe
@@ -449,17 +469,23 @@ export function ListConvosProviderInner({
)
}
} else if (ChatBskyConvoDefs.isLogReadMessage(log)) {
updateConvoInAllLists(log.convoId, convo => ({
...convo,
unreadCount: 0,
rev: log.rev,
}))
updateConvoInAllLists(
log.convoId,
withRevGuard(log.rev, convo => ({
...convo,
unreadCount: 0,
rev: log.rev,
})),
)
} else if (ChatBskyConvoDefs.isLogReadConvo(log)) {
updateConvoInAllLists(log.convoId, convo => ({
...convo,
unreadCount: 0,
rev: log.rev,
}))
updateConvoInAllLists(
log.convoId,
withRevGuard(log.rev, convo => ({
...convo,
unreadCount: 0,
rev: log.rev,
})),
)
} else if (ChatBskyConvoDefs.isLogAcceptConvo(log)) {
const requestQueries =
queryClient.getQueriesData<ConvoListQueryData>({
@@ -472,13 +498,37 @@ export function ListConvosProviderInner({
if (foundConvo) break
}
if (!foundConvo) {
// Use continue (not return) so the remaining logs in this batch
// still apply - the bus advances its cursor past this batch, so a
// dropped log is never redelivered.
debouncedRefetch()
return
continue
}
if (log.rev <= foundConvo.rev) {
continue
}
const acceptedConvo: ChatBskyConvoDefs.ConvoView = {
...foundConvo,
status: 'accepted',
rev: log.rev,
}
// Flip status to 'accepted' in every cache that already holds this
// convo - including 'all'-status caches like the provider's
// always-mounted unread query, which the request->accepted move
// below otherwise never touches. Without this the stale 'all' copy
// keeps status: 'request', and the next isLogCreateMessage can seed
// foundConvo from it and resurrect the convo in the requests inbox.
// Runs before the delete-from-request below: it updates in place
// (never inserts), so the 'request' caches get the accepted copy and
// are then cleared by the delete, leaving no stale request entries.
updateConvoInAllLists(
log.convoId,
withRevGuard(log.rev, convo => ({
...convo,
status: 'accepted',
rev: log.rev,
})),
)
queryClient.setQueriesData(
{queryKey: RQKEY_PARTIAL('request')},
(old?: ConvoListQueryData) => optimisticDelete(log.convoId, old),
@@ -519,60 +569,75 @@ export function ListConvosProviderInner({
},
)
} else if (ChatBskyConvoDefs.isLogMuteConvo(log)) {
mutateConvoView(log.convoId, convo => ({
...convo,
muted: true,
rev: log.rev,
}))
mutateConvoView(
log.convoId,
withRevGuard(log.rev, convo => ({
...convo,
muted: true,
rev: log.rev,
})),
)
} else if (ChatBskyConvoDefs.isLogUnmuteConvo(log)) {
mutateConvoView(log.convoId, convo => ({
...convo,
muted: false,
rev: log.rev,
}))
mutateConvoView(
log.convoId,
withRevGuard(log.rev, convo => ({
...convo,
muted: false,
rev: log.rev,
})),
)
} else if (ChatBskyConvoDefs.isLogLockConvo(log)) {
mutateConvoView(log.convoId, convo => {
if (ChatBskyConvoDefs.isGroupConvo(convo.kind)) {
return {
...convo,
kind: {...convo.kind, lockStatus: 'locked'},
rev: log.rev,
mutateConvoView(
log.convoId,
withRevGuard(log.rev, convo => {
if (ChatBskyConvoDefs.isGroupConvo(convo.kind)) {
return {
...convo,
kind: {...convo.kind, lockStatus: 'locked'},
rev: log.rev,
}
}
}
return {...convo, rev: log.rev}
})
return {...convo, rev: log.rev}
}),
)
// The log event doesn't say whether the lock is forced by a
// moderation override, so refetch to pick up the flag.
void queryClient.invalidateQueries({
queryKey: CONVO_KEY(log.convoId),
})
} else if (ChatBskyConvoDefs.isLogUnlockConvo(log)) {
mutateConvoView(log.convoId, convo => {
if (ChatBskyConvoDefs.isGroupConvo(convo.kind)) {
return {
...convo,
kind: {
...convo.kind,
lockStatus: 'unlocked',
// An unlocked convo cannot be moderation-locked.
lockStatusModerationOverride: false,
},
rev: log.rev,
mutateConvoView(
log.convoId,
withRevGuard(log.rev, convo => {
if (ChatBskyConvoDefs.isGroupConvo(convo.kind)) {
return {
...convo,
kind: {
...convo.kind,
lockStatus: 'unlocked',
// An unlocked convo cannot be moderation-locked.
lockStatusModerationOverride: false,
},
rev: log.rev,
}
}
}
return {...convo, rev: log.rev}
})
return {...convo, rev: log.rev}
}),
)
} else if (ChatBskyConvoDefs.isLogLockConvoPermanently(log)) {
mutateConvoView(log.convoId, convo => {
if (ChatBskyConvoDefs.isGroupConvo(convo.kind)) {
return {
...convo,
kind: {...convo.kind, lockStatus: 'locked-permanently'},
rev: log.rev,
mutateConvoView(
log.convoId,
withRevGuard(log.rev, convo => {
if (ChatBskyConvoDefs.isGroupConvo(convo.kind)) {
return {
...convo,
kind: {...convo.kind, lockStatus: 'locked-permanently'},
rev: log.rev,
}
}
}
return {...convo, rev: log.rev}
})
return {...convo, rev: log.rev}
}),
)
} else if (
ChatBskyConvoDefs.isLogCreateJoinLink(log) ||
ChatBskyConvoDefs.isLogEditJoinLink(log) ||
@@ -592,30 +657,39 @@ export function ListConvosProviderInner({
// Route through mutateConvoView (not updateConvoInAllLists) so the
// single-convo cache updates too, keeping the in-convo requests
// banner in sync.
mutateConvoView(log.convoId, convo =>
applyJoinRequestCountDelta(convo, log.rev, -1),
mutateConvoView(
log.convoId,
withRevGuard(log.rev, convo =>
applyJoinRequestCountDelta(convo, log.rev, -1),
),
)
} else if (ChatBskyConvoDefs.isLogIncomingJoinRequest(log)) {
// Route through mutateConvoView (not updateConvoInAllLists) so the
// single-convo cache updates too, letting the in-convo requests
// banner appear live.
mutateConvoView(log.convoId, convo =>
applyJoinRequestCountDelta(convo, log.rev, 1),
mutateConvoView(
log.convoId,
withRevGuard(log.rev, convo =>
applyJoinRequestCountDelta(convo, log.rev, 1),
),
)
} else if (ChatBskyConvoDefs.isLogReadJoinRequests(log)) {
// The owner marked join requests as read (possibly on another
// device). Zero the unread count but keep the total, mirroring the
// useMarkJoinRequestsRead mutation.
mutateConvoView(log.convoId, convo => {
if (!ChatBskyConvoDefs.isGroupConvo(convo.kind)) {
return {...convo, rev: log.rev}
}
return {
...convo,
kind: {...convo.kind, unreadJoinRequestCount: 0},
rev: log.rev,
}
})
mutateConvoView(
log.convoId,
withRevGuard(log.rev, convo => {
if (!ChatBskyConvoDefs.isGroupConvo(convo.kind)) {
return {...convo, rev: log.rev}
}
return {
...convo,
kind: {...convo.kind, unreadJoinRequestCount: 0},
rev: log.rev,
}
}),
)
} else if (ChatBskyConvoDefs.isLogOutgoingJoinRequest(log)) {
// Viewer isn't in the chat yet, but the inbox surfaces outgoing
// requests, so refetch to pick up the new entry.
@@ -623,8 +697,11 @@ export function ListConvosProviderInner({
} else if (ChatBskyConvoDefs.isLogWithdrawIncomingJoinRequest(log)) {
// A requester rescinded their request to a group the viewer owns.
// Mirror of isLogIncomingJoinRequest: decrement the counts.
mutateConvoView(log.convoId, convo =>
applyJoinRequestCountDelta(convo, log.rev, -1),
mutateConvoView(
log.convoId,
withRevGuard(log.rev, convo =>
applyJoinRequestCountDelta(convo, log.rev, -1),
),
)
} else if (ChatBskyConvoDefs.isLogWithdrawOutgoingJoinRequest(log)) {
// The viewer rescinded their own outgoing join request (possibly on
@@ -634,25 +711,28 @@ export function ListConvosProviderInner({
old => optimisticDeleteJoinRequest(log.convoId, old),
)
} else if (ChatBskyConvoDefs.isLogAddReaction(log)) {
updateConvoInAllLists(log.convoId, convo => {
// add relatedProfiles to members list, but making sure to dedupe
const relatedProfilesSansMembers = (
log.relatedProfiles ?? []
).filter(
profile =>
!convo.members.some(member => member.did === profile.did),
)
return {
...convo,
members: [...convo.members, ...relatedProfilesSansMembers],
lastReaction: {
$type: 'chat.bsky.convo.defs#messageAndReactionView',
reaction: log.reaction,
message: log.message,
},
rev: log.rev,
}
})
updateConvoInAllLists(
log.convoId,
withRevGuard(log.rev, convo => {
// add relatedProfiles to members list, but making sure to dedupe
const relatedProfilesSansMembers = (
log.relatedProfiles ?? []
).filter(
profile =>
!convo.members.some(member => member.did === profile.did),
)
return {
...convo,
members: [...convo.members, ...relatedProfilesSansMembers],
lastReaction: {
$type: 'chat.bsky.convo.defs#messageAndReactionView',
reaction: log.reaction,
message: log.message,
},
rev: log.rev,
}
}),
)
} else if (ChatBskyConvoDefs.isLogAddMember(log)) {
const data = log.message.data
if (
@@ -725,31 +805,35 @@ export function ListConvosProviderInner({
queryClient.setQueriesData(
{queryKey: [RQKEY_ROOT]},
(old?: ConvoListQueryData) =>
optimisticUpdate(log.convoId, old, convo => {
if (
// if the convo is the same
log.convoId === convo.id &&
ChatBskyConvoDefs.isMessageAndReactionView(
convo.lastReaction,
) &&
ChatBskyConvoDefs.isMessageView(log.message) &&
// ...and the message is the same
convo.lastReaction.message.id === log.message.id &&
// ...and the reaction is the same
convo.lastReaction.reaction.sender.did ===
log.reaction.sender.did &&
convo.lastReaction.reaction.value === log.reaction.value
) {
return {
...convo,
// ...remove the reaction. hopefully they didn't react twice in a row!
lastReaction: undefined,
rev: log.rev,
optimisticUpdate(
log.convoId,
old,
withRevGuard(log.rev, convo => {
if (
// if the convo is the same
log.convoId === convo.id &&
ChatBskyConvoDefs.isMessageAndReactionView(
convo.lastReaction,
) &&
ChatBskyConvoDefs.isMessageView(log.message) &&
// ...and the message is the same
convo.lastReaction.message.id === log.message.id &&
// ...and the reaction is the same
convo.lastReaction.reaction.sender.did ===
log.reaction.sender.did &&
convo.lastReaction.reaction.value === log.reaction.value
) {
return {
...convo,
// ...remove the reaction. hopefully they didn't react twice in a row!
lastReaction: undefined,
rev: log.rev,
}
} else {
return convo
}
} else {
return convo
}
}),
}),
),
)
}
}
@@ -905,6 +989,25 @@ export function useOnMarkAsRead() {
)
}
/**
* Wraps a log-driven convo update so it's skipped when the log is not newer
* than the cached convo. Lists are fed by two unsynchronized channels (full
* listConvos refetches and the log stream), so a stale refetch snapshot can be
* written after a log already applied, or a refetch can already include a
* message whose log then arrives and double-counts. Rev comparison as plain
* string comparison is safe here - revs are fixed-width. ConvoView.rev is a
* required field per the lexicon, so convo.rev always exists.
*
* Only used in the log-event paths (which have a log.rev), never in the generic
* optimisticUpdate helper, which mutations without a rev also call.
*/
function withRevGuard(
rev: string,
fn: (convo: ChatBskyConvoDefs.ConvoView) => ChatBskyConvoDefs.ConvoView,
): (convo: ChatBskyConvoDefs.ConvoView) => ChatBskyConvoDefs.ConvoView {
return convo => (rev <= convo.rev ? convo : fn(convo))
}
function optimisticUpdate(
chatId: string,
old?: ConvoListQueryData,