Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
276 changes: 275 additions & 1 deletion src/web/src/hooks/community/community-ws/reconnect-messages.test.ts

Large diffs are not rendered by default.

247 changes: 163 additions & 84 deletions src/web/src/hooks/community/community-ws/reconnect-messages.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import type { InfiniteData, QueryClient, QueryKey } from "@tanstack/react-query"
import type { InfiniteData, Query, QueryClient, QueryKey } from "@tanstack/react-query"
import { apiFetchProfiles, messageProfilePatches } from "@/lib/community/profile-seed"
import { ApiError } from "@/lib/errors"
import { captureChannelMetadataToken, isChannelMetadataTokenCurrent } from "@/hooks/community/channel-metadata"
Expand Down Expand Up @@ -125,7 +125,17 @@ type GapRepairScope = {
serverId?: string
}

const gapRepairs = new WeakMap<QueryClient, Map<string, Promise<void>>>()
type FocusedQueryRepair = { promise: Promise<void>; settled: boolean }

type FocusedMessageRepair = {
token: ReturnType<typeof captureChannelMetadataToken>
queries: Map<Query, FocusedQueryRepair>
promise: Promise<void>
gapPromise: Promise<void>
target: { seq: number }
}

const focusedMessageRepairs = new WeakMap<QueryClient, Map<string, FocusedMessageRepair>>()

function knownFocusedMessageSeq(
queryClient: QueryClient,
Expand Down Expand Up @@ -166,26 +176,7 @@ export function scheduleFocusedMessageGapRepair(
incomingSeq: number,
): Promise<void> | null {
if (incomingSeq <= knownFocusedMessageSeq(queryClient, scope) + 1) return null
let repairs = gapRepairs.get(queryClient)
if (!repairs) {
repairs = new Map()
gapRepairs.set(queryClient, repairs)
}
const key = `${scope.kind}:${scope.scopeId}`
const existing = repairs.get(key)
if (existing) return existing
const repair = reconcileFocusedMessageQueries(
queryClient,
scope.kind,
scope.scopeId,
).catch(() => {
// Realtime delivery is fail-open. A reconnect still runs the same
// authoritative reconciliation path if this best-effort repair fails.
}).finally(() => {
repairs!.delete(key)
})
repairs.set(key, repair)
return repair
return getFocusedMessageRepair(queryClient, scope.kind, scope.scopeId, incomingSeq).gapPromise
}

/**
Expand Down Expand Up @@ -227,6 +218,7 @@ async function fetchCatchUp(
scopeId: string,
cursor: string,
tag: string | null,
target: { seq: number },
): Promise<MessagesPage> {
const messages: Msg[] = []
let latestSeq = 0
Expand All @@ -245,7 +237,10 @@ async function fetchCatchUp(
latestSeq = Math.max(latestSeq, page.latestSeq ?? 0)
hasMoreNewer = page.hasMoreNewer ?? false
newerCursor = page.newerCursor
if (!hasMoreNewer || !newerCursor) break
if (!hasMoreNewer || !newerCursor) {
if (latestSeq < target.seq) continue
break
}
nextCursor = newerCursor
}

Expand Down Expand Up @@ -293,83 +288,167 @@ function mergeReconciledPages(
return { ...cache, pages }
}

export async function reconcileFocusedMessageQueries(
function getFocusedMessageRepair(
queryClient: QueryClient,
kind: "channel" | "dm",
scopeId: string,
): Promise<void> {
incomingSeq = 0,
): FocusedMessageRepair {
const queryKey = kind === "channel"
? communityKeys.channelMessages(scopeId)
: communityKeys.dmMessages(scopeId)
const queries = queryClient.getQueryCache().findAll({
queryKey,
type: "active",
})
let repairs = focusedMessageRepairs.get(queryClient)
if (!repairs) {
repairs = new Map()
focusedMessageRepairs.set(queryClient, repairs)
}
const key = `${kind}:${scopeId}`
const previous = repairs.get(key)
const reusable = previous && isChannelMetadataTokenCurrent(previous.token)
? previous
: undefined
const target = reusable?.target ?? { seq: 0 }
target.seq = Math.max(target.seq, incomingSeq)
if (reusable
&& reusable.queries.size === queries.length
&& queries.every((query) => {
const operation = reusable.queries.get(query)
return operation && !operation.settled
})) return reusable
const token = captureChannelMetadataToken(scopeId)
const operations = queries.map(async (query) => {
const publicationToken = captureCommunityLiveSnapshotToken(queryClient)
const isCurrent = () => isChannelMetadataTokenCurrent(token)
&& queryClient.getQueryCache().find({ queryKey: query.queryKey, exact: true }) === query
// Infinite-query pagination computes its result from the data snapshot at
// fetch start. If that generation completes after reconciliation, TanStack
// can replace the reconciled cache with its stale snapshot. Capture the
// user's pagination intent, cancel that exact generation, then replay the
// same direction against the reconciled cache below.
const pendingDirection = query.state.fetchMeta?.fetchMore?.direction
await queryClient.cancelQueries(
{ queryKey: query.queryKey, exact: true },
{ revert: true, silent: true },
)

if (!isCurrent()) return
let accessDenied = false
try {
const window = warmReconnectWindow(query.queryKey, query.state.data)
if (!window) {
await queryClient.refetchQueries(
{ queryKey: query.queryKey, exact: true, type: "active" },
{ throwOnError: true },
)
return
}
const refreshed = await fetchCurrentWindow(
scopeId,
window.pageParam,
window.tag,
const queryOperations = new Map<Query, FocusedQueryRepair>()
for (const query of queries) {
const existing = reusable?.queries.get(query)
if (existing && !existing.settled) {
queryOperations.set(query, existing)
continue
}
const operation = (async () => {
const publicationToken = captureCommunityLiveSnapshotToken(queryClient)
const isCurrent = () => isChannelMetadataTokenCurrent(token)
&& queryClient.getQueryCache().find({ queryKey: query.queryKey, exact: true }) === query
// Infinite-query pagination computes its result from the data snapshot at
// fetch start. If that generation completes after reconciliation, TanStack
// can replace the reconciled cache with its stale snapshot. Capture the
// user's pagination intent, cancel that exact generation, then replay the
// same direction against the reconciled cache below.
const pendingDirection = query.state.fetchStatus !== "idle"
? query.state.fetchMeta?.fetchMore?.direction
: undefined
await queryClient.cancelQueries(
{ queryKey: query.queryKey, exact: true },
{ revert: true, silent: true },
)

if (!isCurrent()) return
const catchUp = window.cursor !== null
&& (refreshed.latestSeq ?? 0) > window.latestSeq
? await fetchCatchUp(scopeId, window.cursor, window.tag)
: null
if (!isCurrent()) return
publishCommunityMessages(queryClient, {
channelId: scopeId,
messages: catchUp
? [...refreshed.messages, ...catchUp.messages]
: refreshed.messages,
proof: { token: publicationToken, signal: undefined },
})
queryClient.setQueryData<MessageCache>(query.queryKey, (current) => (
isMessageCache(current)
? mergeReconciledPages(current, refreshed, catchUp)
: current
))
} catch (error) {
if (!isCurrent()) return
if (!isDefinitiveAccessDenial(error)) throw error
accessDenied = true
clearDeniedMessageScope(queryClient, kind, scopeId)
} finally {
if (pendingDirection && !accessDenied && isCurrent()) {
await query.fetch(undefined, {
meta: { fetchMore: { direction: pendingDirection } },
let accessDenied = false
let processedTarget = 0
const refresh = async () => {
let window = warmReconnectWindow(query.queryKey, query.state.data)
if (!window) {
await queryClient.refetchQueries(
{ queryKey: query.queryKey, exact: true, type: "active" },
{ throwOnError: true },
)
if (!isCurrent()) return
window = warmReconnectWindow(query.queryKey, query.state.data)
if (!window || target.seq <= window.latestSeq) {
processedTarget = target.seq
return
}
}
let refreshed = await fetchCurrentWindow(
scopeId,
window.pageParam,
window.tag,
)
if (!isCurrent()) return
for (let attempt = 1; window.cursor === null
&& (refreshed.latestSeq ?? 0) < target.seq
&& attempt < MAX_CATCH_UP_PAGES; attempt += 1) {
refreshed = await fetchCurrentWindow(scopeId, window.pageParam, window.tag)
if (!isCurrent()) return
}
const catchUp = window.cursor !== null
&& Math.max(refreshed.latestSeq ?? 0, target.seq) > window.latestSeq
? await fetchCatchUp(scopeId, window.cursor, window.tag, target)
: null
if (!isCurrent()) return
const coveredTarget = target.seq
publishCommunityMessages(queryClient, {
channelId: scopeId,
messages: catchUp
? [...refreshed.messages, ...catchUp.messages]
: refreshed.messages,
proof: { token: publicationToken, signal: undefined },
})
queryClient.setQueryData<MessageCache>(query.queryKey, (current) => (
isMessageCache(current)
? mergeReconciledPages(current, refreshed, catchUp)
: current
))
processedTarget = coveredTarget
}
const reconcile = async () => {
try {
await refresh()
} catch (error) {
if (!isCurrent()) return
if (!isDefinitiveAccessDenial(error)) throw error
accessDenied = true
clearDeniedMessageScope(queryClient, kind, scopeId)
}
}
try {
await reconcile()
} finally {
try {
if (pendingDirection && !accessDenied && isCurrent()) {
await query.fetch(undefined, {
meta: { fetchMore: { direction: pendingDirection } },
})
}
} finally {
if (!accessDenied && isCurrent() && target.seq > processedTarget) {
await reconcile()
}
}
}
})()
const queryRepair = { promise: operation, settled: false }
queryOperations.set(query, queryRepair)
void operation.then(
() => { queryRepair.settled = true },
() => { queryRepair.settled = true },
)
}
const promise = Promise.allSettled([...queryOperations.values()].map((query) => query.promise)).then((settled) => {
if (settled.some((result) => result.status === "rejected")) {
throw new Error("focused messages failed")
}
})
const settled = await Promise.allSettled(operations)
if (settled.some((result) => result.status === "rejected")) {
throw new Error("focused messages failed")
const repair: FocusedMessageRepair = {
token,
queries: queryOperations,
promise,
gapPromise: promise.catch(() => undefined),
target,
}
repairs.set(key, repair)
void promise.finally(() => {
if (repairs.get(key) === repair) repairs.delete(key)
}).catch(() => undefined)
return repair
}

export function reconcileFocusedMessageQueries(
queryClient: QueryClient,
kind: "channel" | "dm",
scopeId: string,
): Promise<void> {
return getFocusedMessageRepair(queryClient, kind, scopeId).promise
}
57 changes: 31 additions & 26 deletions src/web/src/hooks/community/community-ws/reconnect.ts
Original file line number Diff line number Diff line change
Expand Up @@ -124,39 +124,44 @@ async function reconcileCachedServer(queryClient: QueryClient, serverId: string)
}
}

export async function reconcileFocusedCommunityMessages(
queryClient: QueryClient,
sub = useCommunityStore.getState().subscription,
) {
const operations: Promise<unknown>[] = []
if (sub.channelId) {
operations.push(reconcileFocusedMessageQueries(
queryClient,
"channel",
sub.channelId,
))
}
if (sub.secondaryChannelId) {
operations.push(reconcileFocusedMessageQueries(
queryClient,
"channel",
sub.secondaryChannelId,
))
}
if (sub.dmConversationId) {
operations.push(reconcileFocusedMessageQueries(
queryClient,
"dm",
sub.dmConversationId,
))
}
const settled = await Promise.allSettled(operations)
if (settled.some((result) => result.status === "rejected")) throw new Error("focused messages failed")
}

function policyExecutors(
queryClient: QueryClient,
viewerUserId?: string | null,
): Record<CommunityWsReconcilePolicy, () => void | Promise<void>> {
const sub = useCommunityStore.getState().subscription
const queryKeys = queryClient.getQueryCache().getAll().map((query) => query.queryKey)
return {
"focused-messages": async () => {
const operations: Promise<unknown>[] = []
if (sub.channelId) {
operations.push(reconcileFocusedMessageQueries(
queryClient,
"channel",
sub.channelId,
))
}
if (sub.secondaryChannelId) {
operations.push(reconcileFocusedMessageQueries(
queryClient,
"channel",
sub.secondaryChannelId,
))
}
if (sub.dmConversationId) {
operations.push(reconcileFocusedMessageQueries(
queryClient,
"dm",
sub.dmConversationId,
))
}
const settled = await Promise.allSettled(operations)
if (settled.some((result) => result.status === "rejected")) throw new Error("focused messages failed")
},
"focused-messages": () => reconcileFocusedCommunityMessages(queryClient, sub),
"focused-opener": async () => {
const parentMessageId = useCommunityStore.getState().currentChannelMeta?.parentMessageId
if (!parentMessageId) return
Expand Down
Loading
Loading