From e001b8e92b4189286850ea3733d5ffb22b5e3762 Mon Sep 17 00:00:00 2001 From: Gener <39689863+GenerQAQ@users.noreply.github.com> Date: Fri, 2 Oct 2026 00:45:10 +0800 Subject: [PATCH 1/2] fix(web): reconcile foreground messages independently of websocket validation --- .../community-ws/reconnect-messages.test.ts | 168 +++++++++++- .../community-ws/reconnect-messages.ts | 247 ++++++++++++------ .../hooks/community/community-ws/reconnect.ts | 57 ++-- .../use-community-ws-foreground.test.ts | 163 ++++++++++++ .../src/hooks/community/use-community-ws.ts | 12 +- src/web/src/lib/use-user-ws.test.ts | 75 ++++++ src/web/src/lib/use-user-ws.ts | 8 +- 7 files changed, 617 insertions(+), 113 deletions(-) create mode 100644 src/web/src/hooks/community/use-community-ws-foreground.test.ts diff --git a/src/web/src/hooks/community/community-ws/reconnect-messages.test.ts b/src/web/src/hooks/community/community-ws/reconnect-messages.test.ts index 0c8c28824..10e436f8c 100644 --- a/src/web/src/hooks/community/community-ws/reconnect-messages.test.ts +++ b/src/web/src/hooks/community/community-ws/reconnect-messages.test.ts @@ -124,6 +124,96 @@ describe("focused message reconnect catch-up", () => { queryClient.clear() }) + it.each(["channel", "dm"] as const)("shares in-flight %s window work across foreground, reconnect and gap repair", async (kind) => { + const client = new QueryClient() + const key = kind === "channel" ? communityKeys.channelMessages("shared") : communityKeys.dmMessages("shared") + const { unsubscribe } = seedActiveQuery(client, key) + let release!: (page: unknown) => void + apiFetchMock.mockReturnValueOnce(new Promise((resolve) => { release = resolve })) + const cancel = vi.spyOn(client, "cancelQueries") + const first = reconcileFocusedMessageQueries(client, kind, "shared") + await vi.waitFor(() => expect(apiFetchMock).toHaveBeenCalledOnce()) + const reconnect = reconcileFocusedMessageQueries(client, kind, "shared") + const gap = scheduleFocusedMessageGapRepair(client, { kind, scopeId: "shared" }, 5) + expect(reconnect).toBe(first) + expect(cancel).toHaveBeenCalledOnce() + expect(apiFetchMock).toHaveBeenCalledOnce() + apiFetchMock.mockResolvedValueOnce({ messages: [3, 4, 5].map((seq) => ({ id: `m_${seq}`, seq, createdAt: `2026-08-15T00:00:0${seq}.000Z` })), latestSeq: 5, hasMoreNewer: false }) + release({ messages: [], latestSeq: 2, hasMore: false }) + await Promise.all([first, reconnect, gap]) + expect(apiFetchMock).toHaveBeenCalledTimes(2) + expect(client.getQueryData<{ pages: Array<{ messages: Array<{ id: string }> }> }>(key)?.pages.flatMap((page) => page.messages.map((message) => message.id))).toEqual(["m_2", "m_3", "m_4", "m_5", "m_1"]) + apiFetchMock.mockResolvedValueOnce({ messages: [], latestSeq: 5, hasMore: false }) + await reconcileFocusedMessageQueries(client, kind, "shared") + expect(apiFetchMock).toHaveBeenCalledTimes(3) + unsubscribe() + client.clear() + }) + + it("refreshes a completed query variant when another variant still holds the scope owner", async () => { + const client = new QueryClient() + const key = communityKeys.channelMessages("variants") + const taggedKey = [...key, "tag", "selected"] + const subscriptions = [seedActiveQuery(client, key), seedActiveQuery(client, taggedKey)] + let releaseTagged!: (page: unknown) => void + const tagged = new Promise((resolve) => { releaseTagged = resolve }) + const fresh = { messages: [3, 4, 5].map((seq) => ({ id: `m_${seq}`, seq, createdAt: `2026-08-15T00:00:0${seq}.000Z` })), latestSeq: 5, hasMoreNewer: false } + apiFetchMock.mockImplementation((url: string) => { + if (url.includes("since=")) return Promise.resolve(fresh) + if (url.includes("tag=")) return tagged + return Promise.resolve({ messages: [], latestSeq: apiFetchMock.mock.calls.length > 2 ? 5 : 2, hasMore: false }) + }) + const first = reconcileFocusedMessageQueries(client, "channel", "variants") + await vi.waitFor(() => expect(publishCommunityMessagesMock).toHaveBeenCalledOnce()) + const gap = scheduleFocusedMessageGapRepair(client, { kind: "channel", scopeId: "variants" }, 5) + await vi.waitFor(() => expect(publishCommunityMessagesMock).toHaveBeenCalledTimes(2)) + releaseTagged({ messages: [], latestSeq: 2, hasMore: false }) + await Promise.all([first, gap]) + for (const queryKey of [key, taggedKey]) { + expect(client.getQueryData<{ pages: Array<{ messages: Array<{ id: string }> }> }>(queryKey)?.pages.flatMap((page) => page.messages.map((message) => message.id))).toEqual(["m_2", "m_3", "m_4", "m_5", "m_1"]) + } + expect(apiFetchMock.mock.calls.filter(([url]) => String(url).includes("tag=") && !String(url).includes("since="))).toHaveLength(1) + for (const subscription of subscriptions) subscription.unsubscribe() + client.clear() + }) + + it.each(["account", "permission", "replacement"] as const)("does not share an old repair after %s changes or let its cleanup remove the new owner", async (change) => { + const client = new QueryClient() + useCommunityWsStore.getState().activateProfileAccount("a") + const key = communityKeys.channelMessages("changed") + const subscriptions = [seedActiveQuery(client, key).unsubscribe] + let rejectOld!: (reason: unknown) => void + let releaseNew!: (page: unknown) => void + apiFetchMock.mockReturnValueOnce(new Promise((_, reject) => { rejectOld = reject })) + const first = reconcileFocusedMessageQueries(client, "channel", "changed") + await vi.waitFor(() => expect(apiFetchMock).toHaveBeenCalledOnce()) + if (change === "account") useCommunityWsStore.getState().activateProfileAccount("b") + else if (change === "permission") { + useCommunityWsStore.getState().revokeChannelAccess("server", "changed") + useCommunityWsStore.getState().rememberChannelAccess("server", "changed") + } else { + client.removeQueries({ queryKey: key, exact: true }) + subscriptions.push(seedActiveQuery(client, key).unsubscribe) + } + apiFetchMock.mockReturnValueOnce(new Promise((resolve) => { releaseNew = resolve })) + const second = reconcileFocusedMessageQueries(client, "channel", "changed") + expect(second).not.toBe(first) + await vi.waitFor(() => expect(apiFetchMock).toHaveBeenCalledTimes(2)) + rejectOld(new ApiError("old denial", 403)) + await first + expect(client.getQueryData(key)).toBeDefined() + expect(reconcileFocusedMessageQueries(client, "channel", "changed")).toBe(second) + expect(scheduleFocusedMessageGapRepair(client, { kind: "channel", scopeId: "changed" }, 5)).not.toBeNull() + expect(apiFetchMock).toHaveBeenCalledTimes(2) + apiFetchMock.mockResolvedValueOnce({ messages: [], latestSeq: 5, hasMoreNewer: false }) + releaseNew({ messages: [{ id: "current", seq: 2, createdAt: "2026-08-15T00:00:02.000Z" }], latestSeq: 2, hasMore: false }) + await second + expect(publishCommunityMessagesMock).toHaveBeenCalledOnce() + expect(client.getQueryData<{ pages: Array<{ messages: Array<{ id: string }> }> }>(key)?.pages[0].messages[0].id).toBe("current") + for (const unsubscribe of subscriptions) unsubscribe() + client.clear() + }) + it("does not repair exact-next, duplicate, or out-of-order frames", () => { const queryClient = new QueryClient() const queryKey = communityKeys.channelMessages("ch_contiguous") @@ -395,7 +485,12 @@ describe("focused message reconnect catch-up", () => { latestSeq: 3, }) - await reconcileFocusedMessageQueries(queryClient, "channel", "ch_race") + const cancel = vi.spyOn(queryClient, "cancelQueries") + const foreground = reconcileFocusedMessageQueries(queryClient, "channel", "ch_race") + const reconnect = reconcileFocusedMessageQueries(queryClient, "channel", "ch_race") + expect(reconnect).toBe(foreground) + await Promise.all([foreground, reconnect]) + expect(cancel).toHaveBeenCalledOnce() resolveStaleOlder({ messages: [{ id: "stale_m_1" }], hasMore: false, @@ -415,6 +510,77 @@ describe("focused message reconnect catch-up", () => { unsubscribe() }) + it("consumes a new gap arriving while the one pagination replay is still pending", async () => { + const client = new QueryClient() + const key = communityKeys.channelMessages("late-replay") + let olderCalls = 0 + const newest = { messages: [{ id: "m_2", seq: 2, createdAt: "2026-08-15T00:00:02.000Z" }], latestSeq: 2, hasMore: true, cursor: "older" } + const older = { messages: [{ id: "m_1", seq: 1, createdAt: "2026-08-15T00:00:01.000Z" }], latestSeq: 2, hasMore: false } + let releaseStale!: (value: typeof older) => void + let releaseReplay!: (value: typeof older) => void + const observer = new InfiniteQueryObserver(client, { + queryKey: key, + queryFn: async ({ pageParam }) => { + if (pageParam.mode === "newest") return newest + olderCalls += 1 + return new Promise((resolve) => { + if (olderCalls === 1) releaseStale = resolve + else releaseReplay = resolve + }) + }, + initialPageParam: { mode: "newest" }, + getNextPageParam: (last) => last.hasMore ? { mode: "older" } : undefined, + }) + const unsubscribe = observer.subscribe(() => undefined) + await vi.waitFor(() => expect(observer.getCurrentResult().isSuccess).toBe(true)) + const pagination = observer.fetchNextPage() + await vi.waitFor(() => expect(olderCalls).toBe(1)) + apiFetchMock.mockResolvedValueOnce(newest).mockResolvedValueOnce({ ...newest, latestSeq: 5 }).mockResolvedValueOnce({ + messages: [3, 4, 5].map((seq) => ({ id: `m_${seq}`, seq, createdAt: `2026-08-15T00:00:0${seq}.000Z` })), latestSeq: 5, hasMoreNewer: false, + }) + const cancel = vi.spyOn(client, "cancelQueries") + const foreground = reconcileFocusedMessageQueries(client, "channel", "late-replay") + await vi.waitFor(() => expect(olderCalls).toBe(2)) + const gap = scheduleFocusedMessageGapRepair(client, { kind: "channel", scopeId: "late-replay" }, 5) + expect(apiFetchMock).toHaveBeenCalledOnce() + releaseReplay(older) + releaseStale(older) + await Promise.all([foreground, gap, pagination]) + expect(apiFetchMock).toHaveBeenCalledTimes(3) + expect(cancel).toHaveBeenCalledOnce() + expect(olderCalls).toBe(2) + expect(client.getQueryData<{ pages: Array<{ messages: Array<{ id: string }> }> }>(key)?.pages.flatMap((entry) => entry.messages.map((row) => row.id))).toEqual(["m_2", "m_3", "m_4", "m_5", "m_1"]) + unsubscribe() + client.clear() + }) + + it("does not replay completed pagination whose fetch metadata remains idle", async () => { + const client = new QueryClient() + const key = communityKeys.channelMessages("idle-pagination") + let olderCalls = 0 + const newest = { messages: [{ id: "m_2", seq: 2, createdAt: "2026-08-15T00:00:02.000Z" }], latestSeq: 2, hasMore: true, cursor: "older-2" } + const observer = new InfiniteQueryObserver(client, { + queryKey: key, + queryFn: async ({ pageParam }) => { + if (pageParam.mode === "newest") return newest + olderCalls += 1 + return { messages: [{ id: "m_1", seq: 1 }], latestSeq: 2, hasMore: true, cursor: "older-1" } + }, + initialPageParam: { mode: "newest", cursor: "" }, + getNextPageParam: (last) => last.hasMore ? { mode: "older", cursor: last.cursor } : undefined, + }) + const unsubscribe = observer.subscribe(() => undefined) + await vi.waitFor(() => expect(observer.getCurrentResult().isSuccess).toBe(true)) + await observer.fetchNextPage() + expect(client.getQueryState(key)?.fetchStatus).toBe("idle") + expect(client.getQueryState(key)?.fetchMeta?.fetchMore?.direction).toBe("forward") + apiFetchMock.mockResolvedValue(newest) + await reconcileFocusedMessageQueries(client, "channel", "idle-pagination") + expect(olderCalls).toBe(1) + unsubscribe() + client.clear() + }) + it.each([ ["channel", communityKeys.channelMessages("ch_empty"), "ch_empty"], ["dm", communityKeys.dmMessages("dm_empty"), "dm_empty"], diff --git a/src/web/src/hooks/community/community-ws/reconnect-messages.ts b/src/web/src/hooks/community/community-ws/reconnect-messages.ts index 33e8ca70a..10d99eebe 100644 --- a/src/web/src/hooks/community/community-ws/reconnect-messages.ts +++ b/src/web/src/hooks/community/community-ws/reconnect-messages.ts @@ -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" @@ -125,7 +125,17 @@ type GapRepairScope = { serverId?: string } -const gapRepairs = new WeakMap>>() +type FocusedQueryRepair = { promise: Promise; settled: boolean } + +type FocusedMessageRepair = { + token: ReturnType + queries: Map + promise: Promise + gapPromise: Promise + target: { seq: number } +} + +const focusedMessageRepairs = new WeakMap>() function knownFocusedMessageSeq( queryClient: QueryClient, @@ -166,26 +176,7 @@ export function scheduleFocusedMessageGapRepair( incomingSeq: number, ): Promise | 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 } /** @@ -227,6 +218,7 @@ async function fetchCatchUp( scopeId: string, cursor: string, tag: string | null, + target: { seq: number }, ): Promise { const messages: Msg[] = [] let latestSeq = 0 @@ -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 } @@ -293,11 +288,12 @@ function mergeReconciledPages( return { ...cache, pages } } -export async function reconcileFocusedMessageQueries( +function getFocusedMessageRepair( queryClient: QueryClient, kind: "channel" | "dm", scopeId: string, -): Promise { + incomingSeq = 0, +): FocusedMessageRepair { const queryKey = kind === "channel" ? communityKeys.channelMessages(scopeId) : communityKeys.dmMessages(scopeId) @@ -305,71 +301,154 @@ export async function reconcileFocusedMessageQueries( 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() + 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(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(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 { + return getFocusedMessageRepair(queryClient, kind, scopeId).promise } diff --git a/src/web/src/hooks/community/community-ws/reconnect.ts b/src/web/src/hooks/community/community-ws/reconnect.ts index c1c4f5d65..92ebc95e5 100644 --- a/src/web/src/hooks/community/community-ws/reconnect.ts +++ b/src/web/src/hooks/community/community-ws/reconnect.ts @@ -124,6 +124,36 @@ async function reconcileCachedServer(queryClient: QueryClient, serverId: string) } } +export async function reconcileFocusedCommunityMessages( + queryClient: QueryClient, + sub = useCommunityStore.getState().subscription, +) { + const operations: Promise[] = [] + 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, @@ -131,32 +161,7 @@ function policyExecutors( const sub = useCommunityStore.getState().subscription const queryKeys = queryClient.getQueryCache().getAll().map((query) => query.queryKey) return { - "focused-messages": async () => { - const operations: Promise[] = [] - 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 diff --git a/src/web/src/hooks/community/use-community-ws-foreground.test.ts b/src/web/src/hooks/community/use-community-ws-foreground.test.ts new file mode 100644 index 000000000..5c67af44e --- /dev/null +++ b/src/web/src/hooks/community/use-community-ws-foreground.test.ts @@ -0,0 +1,163 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest" +import { QueryObserver } from "@tanstack/react-query" +import { communityKeys } from "@/lib/query-keys" +import { useCommunityStore } from "@/stores/community" +import { useCommunityWsStore } from "@/stores/community/ws" +import { + capturedOnMessage, + capturedOnReconnect, + capturedQueryClient, + capturedUseUserWsOptions, + cleanupCommunityWsHarness, + flushEffects, + getCommunityApiFetchMock, + messageCreate, + mountHook, + resetCommunityWsHarness, + resetHookMemoization, + unmountHook, +} from "./community-ws/test-harness" + +const subscriptions: Array<() => void> = [] +const message = (id: string, seq = 1) => ({ + id, seq, type: "chat" as const, content: id, authorId: "author", + createdAt: `2026-08-15T00:00:0${seq}.000Z`, +}) +const page = (id: string) => ({ messages: [message(id)], latestSeq: 1, hasMore: false }) + +function seed(id: string, kind: "channel" | "dm" = "channel") { + const key = kind === "channel" ? communityKeys.channelMessages(id) : communityKeys.dmMessages(id) + capturedQueryClient.setQueryData(key, { pages: [page(`old-${id}`)], pageParams: [{ mode: "newest" }] }) + const observer = new QueryObserver(capturedQueryClient, { + queryKey: key, + queryFn: async () => ({ pages: [page(`query-${id}`)], pageParams: [{ mode: "newest" }] }), + staleTime: Infinity, + }) + subscriptions.push(observer.subscribe(() => undefined)) + return key +} + +beforeEach(resetCommunityWsHarness) +afterEach(async () => { + for (const unsubscribe of subscriptions.splice(0)) unsubscribe() + await cleanupCommunityWsHarness() +}) + +describe("community foreground message reconciliation", () => { + it.each(["split", "dm"] as const)("starts only focused %s messages while transport access is disconnected", async (layout) => { + const store = useCommunityStore.getState() + const secondary = Symbol("parent") + if (layout === "split") { + store.subscribe({ channelId: "primary" }) + store.claimSecondaryChannel(secondary, "parent") + seed("primary") + seed("parent") + } else { + store.subscribe({ dmConversationId: "dm" }) + seed("dm", "dm") + } + seed("unrelated") + await mountHook() + flushEffects() + useCommunityWsStore.getState().markAccessDisconnected() + getCommunityApiFetchMock().mockImplementation(async (path) => page(String(path))) + const invalidate = vi.spyOn(capturedQueryClient, "invalidateQueries") + await capturedUseUserWsOptions!.onForeground!() + const paths = getCommunityApiFetchMock().mock.calls.map(([path]) => path) + expect(paths).toEqual(layout === "split" ? [ + "/api/community/channels/primary/messages", + "/api/community/channels/parent/messages", + ] : ["/api/community/channels/dm/messages"]) + expect(invalidate).not.toHaveBeenCalled() + if (layout === "split") { + store.releaseSecondaryChannel(secondary) + getCommunityApiFetchMock().mockClear() + await capturedUseUserWsOptions!.onForeground!() + expect(getCommunityApiFetchMock().mock.calls.map(([path]) => path)).toEqual([ + "/api/community/channels/primary/messages", + ]) + } + }) + + it("shares a held foreground window with reconnect and a later gap frame without duplicate cancellation", async () => { + useCommunityStore.getState().subscribe({ channelId: "focused" }) + const key = seed("focused") + await mountHook() + flushEffects() + let release!: (value: unknown) => void + getCommunityApiFetchMock().mockImplementation(async (path) => { + if (path === "/api/community/channels/focused/messages") return new Promise((resolve) => { release = resolve }) + if (String(path).includes("/messages?since=")) return { + messages: [message("missed", 2), message("live", 3)], latestSeq: 3, hasMoreNewer: false, + } + if (path === "/api/community/users/me/read-state") return { revision: 0, readStates: [] } + throw new Error(`unexpected fetch: ${String(path)}`) + }) + const cancel = vi.spyOn(capturedQueryClient, "cancelQueries") + const foreground = capturedUseUserWsOptions!.onForeground!() + await vi.waitFor(() => expect(release).toBeTypeOf("function")) + const duplicate = capturedUseUserWsOptions!.onForeground!() + const reconnect = capturedOnReconnect!({ reconnectDurationMs: 1_000 }) + const event = messageCreate("focused", "live") + event.message.seq = 3 + capturedOnMessage!(event) + await Promise.resolve() + expect(getCommunityApiFetchMock().mock.calls.filter(([path]) => path === "/api/community/channels/focused/messages")).toHaveLength(1) + expect(cancel.mock.calls.filter(([filters]) => JSON.stringify(filters.queryKey) === JSON.stringify(key))).toHaveLength(1) + release(page("snapshot-before-gap")) + await Promise.all([foreground, duplicate, reconnect]) + expect(getCommunityApiFetchMock().mock.calls.filter(([path]) => String(path).includes("/channels/focused/messages"))).toHaveLength(2) + expect(cancel.mock.calls.filter(([filters]) => JSON.stringify(filters.queryKey) === JSON.stringify(key))).toHaveLength(1) + }) + + it("reconciles a new target independently and never applies the previous target response to it", async () => { + useCommunityStore.getState().subscribe({ channelId: "a" }) + seed("a") + const currentKey = seed("b") + await mountHook() + flushEffects() + let releaseOld!: (value: unknown) => void + getCommunityApiFetchMock().mockImplementation(async (path) => { + if (path === "/api/community/channels/a/messages") return new Promise((resolve) => { releaseOld = resolve }) + if (path === "/api/community/channels/b/messages") return page("current-b") + throw new Error(`unexpected fetch: ${String(path)}`) + }) + const old = capturedUseUserWsOptions!.onForeground!() + await vi.waitFor(() => expect(releaseOld).toBeTypeOf("function")) + useCommunityStore.getState().subscribe({ channelId: "b" }) + await capturedUseUserWsOptions!.onForeground!() + releaseOld(page("late-a")) + await old + expect(capturedQueryClient.getQueryData<{ pages: Array<{ messages: Array<{ id: string }> }> }>(currentKey)?.pages.flatMap((entry) => entry.messages.map((row) => row.id))).toEqual(["current-b", "old-b"]) + }) + + it("does not let an old account callback start work after the same hook renders another account", async () => { + useCommunityWsStore.getState().activateProfileAccount("a") + useCommunityStore.getState().subscribe({ channelId: "account-target" }) + seed("account-target") + await mountHook({ viewerUserId: "a" }) + flushEffects() + const old = capturedUseUserWsOptions!.onForeground! + useCommunityWsStore.getState().activateProfileAccount("b") + await old() + expect(getCommunityApiFetchMock()).not.toHaveBeenCalled() + resetHookMemoization() + await mountHook({ viewerUserId: "b" }) + await old() + expect(getCommunityApiFetchMock()).not.toHaveBeenCalled() + getCommunityApiFetchMock().mockResolvedValue(page("current-account")) + await capturedUseUserWsOptions!.onForeground!() + expect(getCommunityApiFetchMock()).toHaveBeenCalledOnce() + }) + + it("does not let a disposed owner start work from an old foreground callback", async () => { + useCommunityStore.getState().subscribe({ channelId: "disposed" }) + seed("disposed") + await mountHook() + flushEffects() + const old = capturedUseUserWsOptions!.onForeground! + unmountHook() + await old() + expect(getCommunityApiFetchMock()).not.toHaveBeenCalled() + }) +}) diff --git a/src/web/src/hooks/community/use-community-ws.ts b/src/web/src/hooks/community/use-community-ws.ts index 30a057ab4..0fa29104f 100644 --- a/src/web/src/hooks/community/use-community-ws.ts +++ b/src/web/src/hooks/community/use-community-ws.ts @@ -12,7 +12,10 @@ import { SEEN_DELIVERY_OPERATION_TRIM_TO, useCommunityWsStore, } from "@/stores/community/ws" -import { reconcileCommunityWsReconnect } from "@/hooks/community/community-ws/reconnect" +import { + reconcileCommunityWsReconnect, + reconcileFocusedCommunityMessages, +} from "@/hooks/community/community-ws/reconnect" import { dispatchCommunityWsEvent, dispatchCommunityWsEvents, @@ -528,10 +531,17 @@ export function useCommunityWs(options?: UseCommunityWsOptions): void { } await reconcileAccountReadState(queryClient, { surfaceMode: "non-inbox" }) }, [queryClient]) + const handleForeground = useCallback(() => { + if (inboxRefreshOwner.current?.disposed + || viewerUserIdRef.current !== viewerUserId + || (viewerUserId !== null && useCommunityWsStore.getState().profileViewerId !== viewerUserId)) return + return reconcileFocusedCommunityMessages(queryClient) + }, [queryClient, viewerUserId]) const { send, reconnectNow } = useUserWs(handleMessage, { onReconnect: handleReconnect, onDisconnect: useCommunityWsStore.getState().markAccessDisconnected, onAuthenticated: handleAuthenticated, + onForeground: handleForeground, onConnectionStateChange: handleConnectionStateChange, requestDaemonStatusOnAuth: false, }) diff --git a/src/web/src/lib/use-user-ws.test.ts b/src/web/src/lib/use-user-ws.test.ts index 1556b14b6..66d9cd8f6 100644 --- a/src/web/src/lib/use-user-ws.test.ts +++ b/src/web/src/lib/use-user-ws.test.ts @@ -1650,6 +1650,81 @@ describe("useUserWs", () => { expect(ws.closed).toBe(false) }) + it("notifies foreground data work before validation and again for changed targets during a pending probe", async () => { + setupTokenFetch() + let finish!: () => void + const dataWork = new Promise((resolve) => { finish = resolve }) + const onForeground = vi.fn(() => dataWork) + const onReconnect = vi.fn() + await mountHook(vi.fn(), { onForeground, onReconnect, requestDaemonStatusOnAuth: false }) + const ws = MockWebSocket.instances[0]! + ws.simulateOpen() + ws.simulateMessage({ type: "auth.ok" }) + + dispatchHiddenToVisible() + expect(onForeground).toHaveBeenCalledOnce() + expect(connectionPings(ws)).toHaveLength(1) + expect(onReconnect).not.toHaveBeenCalled() + dispatchWindowFocus() + mockDocument.dispatch("resume") + dispatchPageShow(true) + expect(onForeground).toHaveBeenCalledTimes(4) + expect(connectionPings(ws)).toHaveLength(1) + + const [{ nonce }] = connectionPings(ws) + ws.simulateMessage({ type: "connection.pong", nonce }) + expect(ws.closed).toBe(false) + expect(MockWebSocket.instances).toEqual([ws]) + expect(onReconnect).not.toHaveBeenCalled() + finish() + await flushPromises() + }) + + it("notifies foreground data work while token or authentication is still pending", async () => { + const token = deferred() + mockFetch.mockReturnValueOnce(token.promise) + const onForeground = vi.fn() + await mountHook(vi.fn(), { onForeground, requestDaemonStatusOnAuth: false }) + dispatchWindowFocus() + expect(onForeground).toHaveBeenCalledOnce() + expect(MockWebSocket.instances).toHaveLength(0) + token.resolve({ ok: true, json: async () => ({ userId: "user-1", token: "test-token" }) } as Response) + await flushPromises() + const ws = MockWebSocket.instances[0]! + dispatchWindowFocus() + expect(onForeground).toHaveBeenCalledTimes(2) + expect(ws.readyState).toBe(MockWebSocket.CONNECTING) + ws.simulateOpen() + dispatchWindowFocus() + expect(onForeground).toHaveBeenCalledTimes(3) + expect(connectionPings(ws)).toHaveLength(0) + }) + + it("keeps foreground data failures independent of socket validation and blocks hidden/offline work", async () => { + setupTokenFetch() + const warn = vi.spyOn(console, "warn").mockImplementation(() => undefined) + const onForeground = vi.fn(() => { throw new Error("data unavailable") }) + await mountHook(vi.fn(), { onForeground, requestDaemonStatusOnAuth: false }) + const ws = MockWebSocket.instances[0]! + ws.simulateOpen() + ws.simulateMessage({ type: "auth.ok" }) + mockDocument.visibilityState = "hidden" + mockDocument.dispatch("visibilitychange") + dispatchWindowFocus() + expect(onForeground).not.toHaveBeenCalled() + mockNavigator.onLine = false + mockWindow.dispatch("offline") + mockDocument.visibilityState = "visible" + mockDocument.dispatch("visibilitychange") + expect(onForeground).not.toHaveBeenCalled() + mockNavigator.onLine = true + mockWindow.dispatch("online") + expect(onForeground).toHaveBeenCalledOnce() + expect(connectionPings(ws)).toHaveLength(1) + expect(warn).toHaveBeenCalledWith("[ws] lifecycle callback threw", { callback: "foreground" }) + warn.mockRestore() + }) + it("coalesces an offline-online signal pair while foreground validation is pending", async () => { setupTokenFetch() await mountHook(vi.fn(), { requestDaemonStatusOnAuth: false }) diff --git a/src/web/src/lib/use-user-ws.ts b/src/web/src/lib/use-user-ws.ts index c2e9c99c9..99740d70d 100644 --- a/src/web/src/lib/use-user-ws.ts +++ b/src/web/src/lib/use-user-ws.ts @@ -127,12 +127,13 @@ export type UseUserWsOptions = { onReconnect?: (info: { reconnectDurationMs: number }) => void | Promise onDisconnect?: () => void | Promise onAuthenticated?: () => void | Promise + onForeground?: () => void | Promise onConnectionStateChange?: (phase: UserWsConnectionPhase) => void | Promise requestDaemonStatusOnAuth?: boolean } function runLifecycleCallback( - name: "authenticated" | "connection-state" | "disconnect" | "reconnect", + name: "authenticated" | "connection-state" | "disconnect" | "foreground" | "reconnect", callback: (() => void | Promise) | undefined, ) { if (!callback) return @@ -175,6 +176,7 @@ export function useUserWs( const onReconnectRef = useRef(options?.onReconnect) const onDisconnectRef = useRef(options?.onDisconnect) const onAuthenticatedRef = useRef(options?.onAuthenticated) + const onForegroundRef = useRef(options?.onForeground) const onConnectionStateChangeRef = useRef(options?.onConnectionStateChange) const lastConnectionPhaseRef = useRef(null) const requestDaemonStatusOnAuthRef = useRef(options?.requestDaemonStatusOnAuth ?? true) @@ -218,9 +220,11 @@ export function useUserWs( useEffect(() => { onDisconnectRef.current = options?.onDisconnect onAuthenticatedRef.current = options?.onAuthenticated + onForegroundRef.current = options?.onForeground onConnectionStateChangeRef.current = options?.onConnectionStateChange }, [ options?.onAuthenticated, + options?.onForeground, options?.onConnectionStateChange, options?.onDisconnect, ]) @@ -921,6 +925,8 @@ export function useUserWs( const recoveryNeeded = forceValidation || connectionValidationNeededRef.current if (!recoveryNeeded) return + runLifecycleCallback("foreground", onForegroundRef.current) + if (isOffline() || isPageHidden()) return if (pendingTokenRef.current) return const ws = wsRef.current From 86f1d345d95db52cb7949730aeb3bd67c83b20a1 Mon Sep 17 00:00:00 2001 From: Gener <39689863+GenerQAQ@users.noreply.github.com> Date: Fri, 2 Oct 2026 01:13:31 +0800 Subject: [PATCH 2/2] test(web): cover foreground reconciliation cache boundaries --- .../community-ws/reconnect-messages.test.ts | 108 ++++++++++++++++++ 1 file changed, 108 insertions(+) diff --git a/src/web/src/hooks/community/community-ws/reconnect-messages.test.ts b/src/web/src/hooks/community/community-ws/reconnect-messages.test.ts index 10e436f8c..e8c453422 100644 --- a/src/web/src/hooks/community/community-ws/reconnect-messages.test.ts +++ b/src/web/src/hooks/community/community-ws/reconnect-messages.test.ts @@ -615,6 +615,114 @@ describe("focused message reconnect catch-up", () => { }, ) + it.each(["channel", "dm"] as const)("rechecks a stale empty %s window when an in-flight foreground repair receives its first gap", async (kind) => { + const client = new QueryClient() + const scopeId = "first-gap" + const key = kind === "channel" ? communityKeys.channelMessages(scopeId) : communityKeys.dmMessages(scopeId) + const { queryFn, unsubscribe } = seedEmptyActiveQuery(client, key) + let release!: (page: unknown) => void + apiFetchMock.mockReturnValueOnce(new Promise((resolve) => { release = resolve })) + const repair = reconcileFocusedMessageQueries(client, kind, scopeId) + await vi.waitFor(() => expect(apiFetchMock).toHaveBeenCalledOnce()) + const gap = scheduleFocusedMessageGapRepair(client, { kind, scopeId }, 2) + expect(gap).not.toBeNull() + apiFetchMock.mockResolvedValueOnce({ + messages: [1, 2].map((seq) => ({ id: `m_${seq}`, seq, createdAt: `2026-08-15T00:00:0${seq}.000Z` })), + latestSeq: 2, + hasMore: false, + }) + release({ messages: [], latestSeq: 0, hasMore: false }) + await Promise.all([repair, gap]) + expect(apiFetchMock).toHaveBeenCalledTimes(2) + expect(apiFetchMock.mock.calls.map(([url]) => url)).toEqual([ + `/api/community/channels/${scopeId}/messages`, + `/api/community/channels/${scopeId}/messages`, + ]) + expect(client.getQueryData<{ pages: Array<{ messages: Array<{ id: string }> }> }>(key)?.pages[0].messages.map((message) => message.id)).toEqual(["m_1", "m_2"]) + expect(queryFn).not.toHaveBeenCalled() + unsubscribe() + client.clear() + }) + + it.each(["channel", "dm"] as const)("bounds stale empty %s reads and allows a later recovery", async (kind) => { + const client = new QueryClient() + const scopeId = "bounded-empty" + const key = kind === "channel" ? communityKeys.channelMessages(scopeId) : communityKeys.dmMessages(scopeId) + const { queryFn, unsubscribe } = seedEmptyActiveQuery(client, key) + apiFetchMock.mockResolvedValue({ messages: [], latestSeq: 0, hasMore: false }) + await scheduleFocusedMessageGapRepair(client, { kind, scopeId }, 2) + expect(apiFetchMock).toHaveBeenCalledTimes(8) + expect(client.getQueryData<{ pages: Array<{ messages: unknown[] }> }>(key)?.pages[0].messages).toEqual([]) + apiFetchMock.mockResolvedValueOnce({ messages: [{ id: "m_later", seq: 1 }], latestSeq: 1, hasMore: false }) + await reconcileFocusedMessageQueries(client, kind, scopeId) + expect(apiFetchMock).toHaveBeenCalledTimes(9) + expect(client.getQueryData<{ pages: Array<{ messages: Array<{ id: string }> }> }>(key)?.pages[0].messages.map((message) => message.id)).toEqual(["m_later"]) + expect(queryFn).not.toHaveBeenCalled() + unsubscribe() + client.clear() + }) + + it.each(["permission", "replacement"] as const)("drops an empty-window retry result after %s changes", async (change) => { + const client = new QueryClient() + const scopeId = "empty-retry-owner" + const key = communityKeys.channelMessages(scopeId) + const subscriptions = [seedEmptyActiveQuery(client, key).unsubscribe] + let release!: (page: unknown) => void + apiFetchMock.mockResolvedValueOnce({ messages: [], latestSeq: 0, hasMore: false }) + apiFetchMock.mockReturnValueOnce(new Promise((resolve) => { release = resolve })) + const gap = scheduleFocusedMessageGapRepair(client, { kind: "channel", scopeId }, 2) + await vi.waitFor(() => expect(apiFetchMock).toHaveBeenCalledTimes(2)) + if (change === "permission") useCommunityWsStore.getState().revokeChannelAccess("server", scopeId) + else { + client.removeQueries({ queryKey: key, exact: true }) + subscriptions.push(seedActiveQuery(client, key).unsubscribe) + } + const current = client.getQueryData(key) + release({ messages: [{ id: "m_obsolete", seq: 2 }], latestSeq: 2, hasMore: false }) + await gap + expect(client.getQueryData(key)).toBe(current) + expect(publishCommunityMessagesMock).not.toHaveBeenCalled() + expect(apiFetchMock).toHaveBeenCalledTimes(2) + for (const unsubscribe of subscriptions) unsubscribe() + client.clear() + }) + + it("preserves a same-query cold reset until its canonical read completes instead of rebuilding pages from old HTTP", async () => { + const client = new QueryClient() + const key = communityKeys.channelMessages("reset-pending") + const oldPage = { messages: [{ id: "m_old", seq: 2 }], latestSeq: 2, hasMore: false } + client.setQueryData(key, { pages: [oldPage], pageParams: [{ mode: "newest" }] }) + let releaseCanonical!: (page: typeof oldPage) => void + const queryFn = vi.fn(() => new Promise((resolve) => { releaseCanonical = resolve })) + const observer = new InfiniteQueryObserver(client, { + queryKey: key, + queryFn, + initialPageParam: { mode: "newest" } as const, + getNextPageParam: () => undefined, + staleTime: Infinity, + }) + const unsubscribe = observer.subscribe(() => undefined) + const query = client.getQueryCache().find({ queryKey: key, exact: true }) + let releaseOld!: (page: typeof oldPage) => void + apiFetchMock.mockReturnValueOnce(new Promise((resolve) => { releaseOld = resolve })) + const repair = reconcileFocusedMessageQueries(client, "channel", "reset-pending") + await vi.waitFor(() => expect(apiFetchMock).toHaveBeenCalledOnce()) + const reset = client.resetQueries({ queryKey: key, exact: true }) + await vi.waitFor(() => expect(queryFn).toHaveBeenCalledOnce()) + expect(client.getQueryCache().find({ queryKey: key, exact: true })).toBe(query) + expect(client.getQueryData(key)).toBeUndefined() + releaseOld(oldPage) + await repair + expect(client.getQueryData(key)).toBeUndefined() + releaseCanonical({ messages: [{ id: "m_current", seq: 3 }], latestSeq: 3, hasMore: false }) + await reset + expect(client.getQueryData<{ pages: Array<{ messages: Array<{ id: string }> }> }>(key)?.pages[0].messages.map((message) => message.id)).toEqual(["m_current"]) + expect(queryFn).toHaveBeenCalledOnce() + expect(apiFetchMock).toHaveBeenCalledOnce() + unsubscribe() + client.clear() + }) + it.each([ ["channel", communityKeys.channelMessages("ch_cold"), "ch_cold"], ["dm", communityKeys.dmMessages("dm_cold"), "dm_cold"],