diff --git a/src/vercel/react/useThreadMessages.ts b/src/vercel/react/useThreadMessages.ts index 6a1b1f14..0b1f0470 100644 --- a/src/vercel/react/useThreadMessages.ts +++ b/src/vercel/react/useThreadMessages.ts @@ -5,7 +5,6 @@ import { type ErrorMessage, type Expand, } from "convex-helpers"; -import { usePaginatedQuery } from "convex-helpers/react"; import { type PaginatedQueryArgs, type UsePaginatedQueryResult, @@ -28,6 +27,7 @@ import type { } from "../../validators.js"; import type { StreamQueryArgs, StreamQuery } from "./types.js"; import { useStreamingUIMessages } from "./useStreamingUIMessages.js"; +import { usePaginatedMessages } from "./useUIMessages.js"; export type MessageDocLike = { order: number; @@ -120,6 +120,11 @@ export function useThreadMessages>( query: Query, args: ThreadMessagesArgs | "skip", options: { + /** + * Minimum page size. The hook may load additional records, and briefly + * report "LoadingMore", when pagination splits a message across canonical + * rows sharing one `order`. + */ initialNumItems: number; stream?: Query extends StreamQuery ? boolean @@ -129,10 +134,10 @@ export function useThreadMessages>( ThreadMessagesResult & { streaming: boolean; key: string } > { // These are full messages - const paginated = usePaginatedQuery( + const paginated = usePaginatedMessages( query, args as PaginatedQueryArgs | "skip", - { initialNumItems: options.initialNumItems }, + options.initialNumItems, ); let startOrder = paginated.results.at(-1)?.order ?? 0; diff --git a/src/vercel/react/useUIMessages.test.ts b/src/vercel/react/useUIMessages.test.ts index b6ef9ba2..dee71843 100644 --- a/src/vercel/react/useUIMessages.test.ts +++ b/src/vercel/react/useUIMessages.test.ts @@ -1,7 +1,10 @@ import { describe, it, expect } from "vitest"; +import { toUIMessages } from "../UIMessages.js"; +import type { MessageDoc } from "../../validators.js"; import { dedupeMessages, mergeUIMessages, + planOldestOrderCompletion, type UIMessageLike, } from "./useUIMessages.js"; @@ -45,6 +48,86 @@ function testUIMessage({ }; } +describe("paginated UI message boundaries", () => { + const partialOrder = [ + { order: 4, stepOrder: 2 }, + { order: 5, stepOrder: 0 }, + ]; + + it("tracks one incomplete order across follow-up pages", () => { + expect( + planOldestOrderCompletion(partialOrder, "CanLoadMore", "ready"), + ).toEqual({ order: 4, loadMore: true }); + expect( + planOldestOrderCompletion( + [...partialOrder, { order: 4, stepOrder: 1 }], + "CanLoadMore", + 4, + ), + ).toEqual({ order: 4, loadMore: true }); + }); + + it("stops after completing the target instead of chasing a new boundary", () => { + expect( + planOldestOrderCompletion( + [ + { order: 4, stepOrder: 0 }, + { order: 4, stepOrder: 2 }, + { order: 3, stepOrder: 2 }, + ], + "CanLoadMore", + 4, + ), + ).toEqual({ order: "idle", loadMore: false }); + }); + + it("stops if a custom query filters out the target's step zero", () => { + expect( + planOldestOrderCompletion( + [...partialOrder, { order: 3, stepOrder: 2 }], + "CanLoadMore", + 4, + ), + ).toEqual({ order: "idle", loadMore: false }); + }); + + it("uses loading status to distinguish automatic and caller pages", () => { + expect(planOldestOrderCompletion(partialOrder, "LoadingMore", 4)).toEqual({ + order: 4, + loadMore: false, + }); + expect( + planOldestOrderCompletion(partialOrder, "LoadingMore", "idle"), + ).toEqual({ order: "ready", loadMore: false }); + }); + + it("resets with a new first page and ignores an already complete boundary", () => { + expect( + planOldestOrderCompletion(partialOrder, "LoadingFirstPage", 4), + ).toEqual({ order: "ready", loadMore: false }); + expect( + planOldestOrderCompletion( + [ + { order: 4, stepOrder: 0 }, + { order: 5, stepOrder: 0 }, + ], + "CanLoadMore", + "ready", + ), + ).toEqual({ order: "idle", loadMore: false }); + expect(planOldestOrderCompletion(partialOrder, "Exhausted", 4)).toEqual({ + order: "idle", + loadMore: false, + }); + }); + + it("stops if realtime updates remove the target order", () => { + expect( + planOldestOrderCompletion([{ order: 5, stepOrder: 0 }], "CanLoadMore", 4), + ).toEqual({ order: "idle", loadMore: false }); + }); +}); + describe("dedupeMessages", () => { it("should prefer messages from messages list when streaming messages are absent", () => { const messages: TestMessage[] = [ @@ -332,3 +415,97 @@ describe("mergeUIMessages", () => { ]); }); }); + +describe("split pagination boundary (issue #193)", () => { + function doc(overrides: Partial): MessageDoc { + return { + _id: `m${overrides.order}-${overrides.stepOrder}`, + _creationTime: 0, + order: 0, + stepOrder: 0, + status: "success", + threadId: "t1", + tool: false, + ...overrides, + }; + } + + // Page starts partway through order 4: the tool call at stepOrder 0 is missing. + const splitPage = [ + doc({ + order: 4, + stepOrder: 1, + tool: true, + message: { + role: "tool", + content: [ + { + type: "tool-result", + toolCallId: "call1", + toolName: "myTool", + output: { type: "text", value: "42" }, + }, + ], + }, + }), + doc({ + order: 4, + stepOrder: 2, + message: { role: "assistant", content: "The answer is 42." }, + text: "The answer is 42.", + }), + ]; + + const anchorPage = [ + doc({ + order: 4, + stepOrder: 0, + tool: true, + message: { + role: "assistant", + content: [ + { + type: "tool-call", + toolCallId: "call1", + toolName: "myTool", + input: "question", + args: "question", + }, + ], + }, + }), + ]; + + it("loads follow-up pages until the split order is whole", () => { + let loaded = splitPage; + let orderToComplete: "ready" | number | "idle" = "ready"; + let loads = 0; + for (let i = 0; i < 5; i++) { + const decision = planOldestOrderCompletion( + loaded, + "CanLoadMore", + orderToComplete, + ); + orderToComplete = decision.order; + if (!decision.loadMore) break; + loads++; + loaded = [...anchorPage, ...loaded]; + } + + expect(loads).toBe(1); + expect(orderToComplete).toBe("idle"); + + const order4 = toUIMessages(loaded).filter((m) => m.order === 4); + expect(order4).toHaveLength(1); + expect(order4[0].stepOrder).toBe(0); + expect(order4[0].key).toBe("t1-4-0"); + expect( + order4[0].parts.filter((p) => p.type === "tool-myTool"), + ).toHaveLength(1); + }); + + it("keys a split order off the wrong row, which is why streaming never settles", () => { + const order4 = toUIMessages(splitPage).filter((m) => m.order === 4); + expect(order4[0].key).not.toBe("t1-4-0"); + }); +}); diff --git a/src/vercel/react/useUIMessages.ts b/src/vercel/react/useUIMessages.ts index 0396adc7..66aeaa88 100644 --- a/src/vercel/react/useUIMessages.ts +++ b/src/vercel/react/useUIMessages.ts @@ -9,13 +9,14 @@ import { type PaginatedQueryArgs, type UsePaginatedQueryResult, } from "convex/react"; -import type { - FunctionArgs, - FunctionReference, - PaginationOptions, - PaginationResult, +import { + getFunctionName, + type FunctionArgs, + type FunctionReference, + type PaginationOptions, + type PaginationResult, } from "convex/server"; -import { useMemo } from "react"; +import { useEffect, useMemo, useRef } from "react"; import type { SyncStreamsReturnValue } from "../client/types.js"; import type { StreamArgs } from "../../validators.js"; import type { StreamQuery } from "./types.js"; @@ -35,6 +36,8 @@ export type UIMessageLike = { role: UIMessage["role"]; }; +type OldestOrderCompletion = "ready" | number | "idle"; + export type UIMessagesQuery< Args = unknown, M extends UIMessageLike = UIMessageLike, @@ -126,6 +129,10 @@ export function useUIMessages>( query: Query, args: UIMessagesQueryArgs | "skip", options: { + /** + * Minimum page size. The hook may load additional records when pagination + * splits the oldest UI message across canonical message rows. + */ initialNumItems: number; stream?: Query extends StreamQuery ? boolean @@ -134,10 +141,10 @@ export function useUIMessages>( }, ): UsePaginatedQueryResult> { // These are full messages - const paginated = usePaginatedQuery( + const paginated = usePaginatedMessages( query, args as PaginatedQueryArgs | "skip", - { initialNumItems: options.initialNumItems }, + options.initialNumItems, ); const startOrder = paginated.results.length @@ -164,6 +171,119 @@ export function useUIMessages>( return merged as UIMessagesQueryResult; } +/** + * usePaginatedQuery, extended to keep the oldest message whole. `initialNumItems` + * is a minimum, not an exact page size. A custom query that filters out a row at + * `stepOrder: 0` while keeping later steps of that order cannot be completed, and + * its oldest message is left split. + */ +export function usePaginatedMessages< + M extends { order: number; stepOrder: number }, + Query extends FunctionReference< + "query", + "public", + any, + PaginationResult & { streams?: SyncStreamsReturnValue } + >, +>(query: Query, args: PaginatedQueryArgs | "skip", numItems: number) { + const paginated = usePaginatedQuery(query, args, { + initialNumItems: numItems, + }); + const orderToComplete = useRef("ready"); + // A new pagination session can skip "LoadingFirstPage" when Convex already + // has the query subscribed, so track identity rather than trusting status. + const queryKey = + args === "skip" + ? "skip" + : `${getFunctionName(query)}:${JSON.stringify(args)}`; + const activeQuery = useRef(queryKey); + + useEffect(() => { + if (activeQuery.current !== queryKey) { + activeQuery.current = queryKey; + orderToComplete.current = "ready"; + } + const decision = planOldestOrderCompletion( + paginated.results, + paginated.status, + orderToComplete.current, + ); + orderToComplete.current = decision.order; + if (decision.loadMore) { + paginated.loadMore(numItems); + } + }, [ + queryKey, + paginated.loadMore, + paginated.results, + paginated.status, + numItems, + ]); + + return paginated; +} + +/** + * UI messages are assembled from every canonical record at the same order. + * If row pagination starts partway through the oldest loaded order, fetch + * follow-up pages until that specific UI message is complete. + */ +export function planOldestOrderCompletion( + messages: Pick[], + status: UsePaginatedQueryResult["status"], + orderToComplete: OldestOrderCompletion, +): { order: OldestOrderCompletion; loadMore: boolean } { + if (status === "LoadingFirstPage") { + return { order: "ready", loadMore: false }; + } + if (status === "LoadingMore") { + // A load started while idle came from the caller. Rearm completion for the + // new boundary. Automatic loads already carry their target order. + return { + order: orderToComplete === "idle" ? "ready" : orderToComplete, + loadMore: false, + }; + } + if (status === "Exhausted") { + return { order: "idle", loadMore: false }; + } + if (status !== "CanLoadMore" || messages.length === 0) { + return { order: orderToComplete, loadMore: false }; + } + + const oldestOrder = Math.min(...messages.map((message) => message.order)); + if (orderToComplete === "idle") { + return { order: "idle", loadMore: false }; + } + + let targetOrder = orderToComplete; + if (targetOrder === "ready") { + const oldestOrderIsComplete = messages.some( + (message) => message.order === oldestOrder && message.stepOrder === 0, + ); + if (oldestOrderIsComplete) { + return { order: "idle", loadMore: false }; + } + targetOrder = oldestOrder; + } + + const targetMessages = messages.filter( + (message) => message.order === targetOrder, + ); + if (targetMessages.length === 0) { + return { order: "idle", loadMore: false }; + } + if (targetMessages.some((message) => message.stepOrder === 0)) { + return { order: "idle", loadMore: false }; + } + if (oldestOrder < targetOrder) { + // The query moved past the target without exposing its step zero. This can + // happen when a custom query filters the anchor; do not chase older orders. + return { order: "idle", loadMore: false }; + } + return { order: targetOrder, loadMore: true }; +} + export function mergeUIMessages( messages: M[], streamMessages: M[],