diff --git a/.changeset/perf-mobile-streaming.md b/.changeset/perf-mobile-streaming.md new file mode 100644 index 000000000000..fea8a0c55988 --- /dev/null +++ b/.changeset/perf-mobile-streaming.md @@ -0,0 +1,5 @@ +--- +"@t3tools/web": patch +--- + +Smoother streaming in chats: bot replies update the conversation once per batch instead of once per token, and typing a message no longer re-renders the whole chat. diff --git a/apps/mobile/src/features/home/HomeScreen.tsx b/apps/mobile/src/features/home/HomeScreen.tsx index 7b7ba69832e7..d6458dc6bbfe 100644 --- a/apps/mobile/src/features/home/HomeScreen.tsx +++ b/apps/mobile/src/features/home/HomeScreen.tsx @@ -45,6 +45,7 @@ import { import { ThreadListV2PendingRow, ThreadListV2Row, + threadListV2TimeLabel, ThreadListV2SettledShelfHeader, ThreadListV2SnoozedShelfHeader, } from "../threads/thread-list-v2-items"; @@ -775,7 +776,10 @@ export function HomeScreen(props: HomeScreenProps) { variant={item.item.variant} snoozed={item.item.snoozed} pinned={item.item.pinned} - snoozePresetMinute={nowMinute} + timeLabel={threadListV2TimeLabel(thread)} + snoozePresetMinute={ + !item.item.snoozed && snoozeEnvironmentIds.has(thread.environmentId) ? nowMinute : null + } snoozeWakeLabelText={item.snoozeWakeLabelText} showTrailingDivider={showTrailingDivider} project={ diff --git a/apps/mobile/src/features/threads/ThreadDetailScreen.tsx b/apps/mobile/src/features/threads/ThreadDetailScreen.tsx index 971e983ee6bc..ecdf662a65c6 100644 --- a/apps/mobile/src/features/threads/ThreadDetailScreen.tsx +++ b/apps/mobile/src/features/threads/ThreadDetailScreen.tsx @@ -57,7 +57,7 @@ import { ControlPill } from "../../components/ControlPill"; import { useAppearancePreferences } from "../settings/appearance/AppearancePreferencesProvider"; import type { ComposerEditorHandle } from "../../components/ComposerEditor"; import type { StatusTone } from "../../components/StatusPill"; -import type { DraftComposerImageAttachment } from "../../lib/composerImages"; +import { useThreadDraftForThread } from "../../state/use-thread-composer-state"; import { CHAT_CONTENT_MAX_WIDTH, type LayoutVariant } from "../../lib/layout"; import { IOS_NAV_BAR_HEIGHT } from "../../lib/layoutMetrics"; import { scopedThreadKey } from "../../lib/scopedEntities"; @@ -78,6 +78,7 @@ import { COMPOSER_COLLAPSED_CHROME, COMPOSER_EXPANDED_CHROME, ThreadComposer, + type ThreadComposerProps, } from "./ThreadComposer"; import { ThreadFeed } from "./ThreadFeed"; import type { ThreadContentPresentation } from "./threadContentPresentation"; @@ -97,8 +98,6 @@ export interface ThreadDetailScreenProps { readonly activePendingUserInputDrafts: Record; readonly activePendingUserInputAnswers: Record> | null; readonly respondingUserInputId: ApprovalRequestId | null; - readonly draftMessage: string; - readonly draftAttachments: ReadonlyArray; readonly connectionStateLabel: EnvironmentConnectionPhase; /** Message sync status for the selected thread (drives the composer status pill). */ readonly threadSyncStatus?: EnvironmentThreadStatus; @@ -149,11 +148,13 @@ function latestStreamingAssistantMessage( ): { readonly id: string; readonly textLength: number } | null { for (let index = feed.length - 1; index >= 0; index -= 1) { const entry = feed[index]; - if (entry?.type !== "message") { + if (entry?.type !== "message" || entry.message.role !== "assistant") { continue; } - if (entry.message.role !== "assistant" || !entry.message.streaming) { - continue; + // Only the newest assistant message can be streaming, so stop there + // instead of walking the whole history after a turn settles. + if (!entry.message.streaming) { + return null; } return { id: entry.message.id, @@ -164,6 +165,20 @@ function latestStreamingAssistantMessage( return null; } +/** Submitted messages land at the tail, so search newest-first. */ +function feedHasMessageNearEnd( + feed: ReadonlyArray, + messageId: MessageId, +): boolean { + for (let index = feed.length - 1; index >= 0; index -= 1) { + const entry = feed[index]; + if (entry?.type === "message" && entry.id === messageId) { + return true; + } + } + return false; +} + function useStreamingHaptics(threadId: ThreadId, feed: ReadonlyArray) { const lastStreamingAssistantRef = useRef<{ readonly id: string; @@ -219,6 +234,29 @@ const USER_INPUT_TOGGLE_TIMING = { easing: Easing.out(Easing.cubic), }; +/** + * Reads the chat's draft next to the composer so a keystroke re-renders the + * composer only, not the thread screen and its feed. + */ +const ThreadDraftComposer = memo(function ThreadDraftComposer( + props: Omit & { + readonly threadId: ThreadId; + }, +) { + const { threadId, ...composerProps } = props; + const { draftMessage, draftAttachments } = useThreadDraftForThread({ + environmentId: props.environmentId, + threadId, + }); + return ( + + ); +}); + export const ThreadDetailScreen = memo(function ThreadDetailScreen(props: ThreadDetailScreenProps) { const insets = useSafeAreaInsets(); const isKeyboardVisible = useKeyboardState((state) => state.isVisible); @@ -476,9 +514,7 @@ export const ThreadDetailScreen = memo(function ThreadDetailScreen(props: Thread submittedMessageId === null || lastScrolledSubmittedMessageIdRef.current === submittedMessageId || contentPresentationKind !== "ready" || - !selectedThreadFeed.some( - (entry) => entry.type === "message" && entry.id === submittedMessageId, - ) + !feedHasMessageNearEnd(selectedThreadFeed, submittedMessageId) ) { return; } @@ -525,9 +561,16 @@ export const ThreadDetailScreen = memo(function ThreadDetailScreen(props: Thread selectedThreadKey, ]); + // Read through a ref so streaming feed updates do not change the send + // callback and re-render the memoized composer on every delta. + const selectedThreadFeedRef = useRef(selectedThreadFeed); + useLayoutEffect(() => { + selectedThreadFeedRef.current = selectedThreadFeed; + }, [selectedThreadFeed]); + const handleSendMessage = useCallback(async () => { const targetThreadKey = selectedThreadKey; - const hasUserMessage = selectedThreadFeed.some( + const hasUserMessage = selectedThreadFeedRef.current.some( (entry) => entry.type === "message" && entry.message.role === "user", ); const messageId = await props.onSendMessage(); @@ -552,7 +595,6 @@ export const ThreadDetailScreen = memo(function ThreadDetailScreen(props: Thread props.onSendMessage, props.selectedThread.latestTurn, props.selectedThreadQueueCount, - selectedThreadFeed, selectedThreadKey, ]); @@ -764,10 +806,9 @@ export const ThreadDetailScreen = memo(function ThreadDetailScreen(props: Thread {/* Hidden (not unmounted) while a user-input request owns the composer slot, so composer drafts and editor state survive. */} - entry.id; +const threadFeedItemType = (entry: ThreadFeedEntry) => + entry.type === "message" ? `message:${entry.message.role}` : entry.type; + +type PlaybackMessage = Extract["message"]; + +function sameMessages( + left: ReadonlyArray, + right: ReadonlyArray, +): boolean { + return left.length === right.length && left.every((message, index) => message === right[index]); +} + +function sameIds(left: ReadonlySet, right: ReadonlySet): boolean { + if (left.size !== right.size) return false; + for (const id of left) { + if (!right.has(id)) return false; + } + return true; +} + export const ThreadFeed = memo(function ThreadFeed(props: ThreadFeedProps) { const navigation = useNavigation(); const replyPlayback = useOptionalReplyPlayback(); const environment = useEnvironmentPresentation(props.environmentId); - const playbackMessages = useMemo( - () => props.feed.flatMap((entry) => (entry.type === "message" ? [entry.message] : [])), - [props.feed], - ); + // Activity-only feed updates keep the previous array so reply playback does not re-observe. + const playbackMessagesRef = useRef>([]); + const playbackMessages = useMemo(() => { + const next = props.feed.flatMap((entry) => (entry.type === "message" ? [entry.message] : [])); + if (sameMessages(playbackMessagesRef.current, next)) return playbackMessagesRef.current; + playbackMessagesRef.current = next; + return next; + }, [props.feed]); useReplyPlaybackThread({ environmentId: props.environmentId, threadId: props.threadId, @@ -1763,28 +1789,6 @@ export const ThreadFeed = memo(function ThreadFeed(props: ThreadFeedProps) { ); const markdownStyles = useMarkdownStyles(onMarkdownLinkPress, renderMarkdownImage); const reviewCommentColors = useReviewCommentColors(); - // LegendList does not invalidate visible rows when only the renderItem closure changes. - // Keep row-local interaction props in extraData so disclosures and copy feedback repaint. - const listAppearanceData = useMemo( - () => ({ - copiedRowId, - expandedWorkRows, - iconSubtleColor, - markdownStyles, - reviewCommentColors, - userBubbleColor, - viewportWidth, - }), - [ - copiedRowId, - expandedWorkRows, - iconSubtleColor, - markdownStyles, - reviewCommentColors, - userBubbleColor, - viewportWidth, - ], - ); const reportHeaderMaterialVisibility = useCallback( (visible: boolean) => { if (headerMaterialVisibleRef.current === visible) { @@ -1953,6 +1957,9 @@ export const ThreadFeed = memo(function ThreadFeed(props: ThreadFeedProps) { ), [presentedFeed, props.anchorMessageId, anchorTopInset], ); + // Kept identity-stable while streaming: it is row render context, and a new Set would + // repaint every visible row on each delta. + const terminalAssistantMessageIdsRef = useRef>(new Set()); const terminalAssistantMessageIds = useMemo(() => { const terminalIdsByTurn = new Map(); for (const entry of props.feed) { @@ -1960,7 +1967,12 @@ export const ThreadFeed = memo(function ThreadFeed(props: ThreadFeedProps) { terminalIdsByTurn.set(entry.message.turnId, entry.message.id); } } - return new Set(terminalIdsByTurn.values()); + const next = new Set(terminalIdsByTurn.values()); + if (sameIds(terminalAssistantMessageIdsRef.current, next)) { + return terminalAssistantMessageIdsRef.current; + } + terminalAssistantMessageIdsRef.current = next; + return next; }, [props.feed]); const unsettledTurnId = props.latestTurn && @@ -2138,30 +2150,31 @@ export const ThreadFeed = memo(function ThreadFeed(props: ThreadFeedProps) { [expandedWorkRows, workingRowHeight, appearance.baseFontSize], ); - const renderItem = useCallback( - (info: { item: ThreadFeedEntry; index: number }) => - renderFeedEntry(info, { - environmentId: props.environmentId, - copiedRowId, - expandedWorkRows, - terminalAssistantMessageIds, - unsettledTurnId, - onCopyWorkRow, - onToggleWorkGroup, - onToggleWorkRow, - onToggleTurnFold, - onPressImage, - onMarkdownLinkPress, - renderMarkdownImage, - iconSubtleColor, - userBubbleColor, - markdownStyles, - reviewCommentColors, - reviewCommentBubbleWidth, - userBubbleMaxWidth, - skills: props.skills, - replyPlayback, - }), + // LegendList repaints visible rows only when their item or extraData changes, never for a + // new renderItem closure. Everything rows read lives here and doubles as extraData. + const feedRenderContext = useMemo( + () => ({ + environmentId: props.environmentId, + copiedRowId, + expandedWorkRows, + terminalAssistantMessageIds, + unsettledTurnId, + onCopyWorkRow, + onToggleWorkGroup, + onToggleWorkRow, + onToggleTurnFold, + onPressImage, + onMarkdownLinkPress, + renderMarkdownImage, + iconSubtleColor, + userBubbleColor, + markdownStyles, + reviewCommentColors, + reviewCommentBubbleWidth, + userBubbleMaxWidth, + skills: props.skills, + replyPlayback, + }), [ copiedRowId, expandedWorkRows, @@ -2185,6 +2198,35 @@ export const ThreadFeed = memo(function ThreadFeed(props: ThreadFeedProps) { replyPlayback, ], ); + const renderItem = useCallback( + (info: { item: ThreadFeedEntry; index: number }) => renderFeedEntry(info, feedRenderContext), + [feedRenderContext], + ); + const loadEarlier = props.loadEarlier; + const listHeader = useMemo( + () => ( + <> + {usesNativeAutomaticInsets ? null : } + + {loadEarlier != null ? ( + + + {loadEarlier.loading ? "Loading earlier turns…" : "Load earlier turns"} + + + ) : null} + + ), + [loadEarlier, props.botId, props.environmentId, topContentInset, usesNativeAutomaticInsets], + ); + const listContentContainerStyle = useMemo( + () => ({ paddingTop: 12, paddingHorizontal: contentHorizontalPadding }), + [contentHorizontalPadding], + ); if (props.contentPresentation.kind === "unavailable") { return ( @@ -2279,12 +2321,11 @@ export const ThreadFeed = memo(function ThreadFeed(props: ThreadFeedProps) { } maintainVisibleContentPosition={maintainVisibleContentPosition} data={presentedFeed} - extraData={listAppearanceData} + extraData={feedRenderContext} renderItem={renderItem} - keyExtractor={(entry) => entry.id} - getItemType={(entry) => - entry.type === "message" ? `message:${entry.message.role}` : entry.type - } + keyExtractor={threadFeedKeyExtractor} + getItemType={threadFeedItemType} + itemsAreEqual={threadFeedEntriesEqual} getFixedItemSize={getFixedItemSize} // Measure rows well before they scroll into view so estimate→actual // corrections land offscreen instead of under the user's finger. @@ -2320,27 +2361,8 @@ export const ThreadFeed = memo(function ThreadFeed(props: ThreadFeedProps) { onMomentumScrollBegin={handleMomentumScrollBegin} onMomentumScrollEnd={handleMomentumScrollEnd} scrollEventThrottle={16} - ListHeaderComponent={ - <> - {usesNativeAutomaticInsets ? null : } - - {props.loadEarlier != null ? ( - - - {props.loadEarlier.loading ? "Loading earlier turns…" : "Load earlier turns"} - - - ) : null} - - } - contentContainerStyle={{ - paddingTop: 12, - paddingHorizontal: contentHorizontalPadding, - }} + ListHeaderComponent={listHeader} + contentContainerStyle={listContentContainerStyle} /> {props.feed.length === 0 && diff --git a/apps/mobile/src/features/threads/ThreadNavigationSidebar.tsx b/apps/mobile/src/features/threads/ThreadNavigationSidebar.tsx index ffabf5b7a845..fddf40f340aa 100644 --- a/apps/mobile/src/features/threads/ThreadNavigationSidebar.tsx +++ b/apps/mobile/src/features/threads/ThreadNavigationSidebar.tsx @@ -73,6 +73,7 @@ import { import { ThreadListV2PendingRow, ThreadListV2Row, + threadListV2TimeLabel, ThreadListV2SettledShelfHeader, ThreadListV2SnoozedShelfHeader, } from "./thread-list-v2-items"; @@ -850,7 +851,12 @@ function ThreadNavigationSidebarPane( variant={item.item.variant} snoozed={item.item.snoozed} pinned={item.item.pinned} - snoozePresetMinute={nowMinute} + timeLabel={threadListV2TimeLabel(thread)} + snoozePresetMinute={ + !item.item.snoozed && snoozeEnvironmentIds.has(thread.environmentId) + ? nowMinute + : null + } snoozeWakeLabelText={item.snoozeWakeLabelText} project={projectByKey.get(scopeKey) ?? null} projectTitle={projectTitleByProjectKey.get(scopeKey)} diff --git a/apps/mobile/src/features/threads/ThreadRouteScreen.tsx b/apps/mobile/src/features/threads/ThreadRouteScreen.tsx index 62de4d605a05..3431211e0db8 100644 --- a/apps/mobile/src/features/threads/ThreadRouteScreen.tsx +++ b/apps/mobile/src/features/threads/ThreadRouteScreen.tsx @@ -1,4 +1,7 @@ -import { NativeStackScreenOptions } from "../../native/StackHeader"; +import { + NativeStackScreenOptions, + type AppNativeStackNavigationOptions, +} from "../../native/StackHeader"; import { StackActions, useFocusEffect, @@ -752,6 +755,55 @@ function ThreadRouteContent( [navigation], ); + const selectedThreadTitle = selectedThread?.title ?? ""; + // Memoized so a composer keystroke does not re-sign the whole header config. + const stackScreenOptions = useMemo( + () => ({ + // Android draws its own in-flow header (AndroidScreenHeader below); + // the native stack header stays iOS-only. + headerShown: Platform.OS !== "android", + headerTitle: selectedThreadTitle, + headerTitleStyle: usesNativeHeaderGlass + ? { + fontSize: 17, + fontWeight: "800", + } + : undefined, + title: selectedThreadTitle, + headerBackVisible: !layout.usesSplitView, + // Compact uses the NATIVE back button when a previous route exists; + // deep links / cold starts get an explicit Home button instead. + // Split view always uses its custom left items. + unstable_headerLeftItems: + Platform.OS === "ios" + ? layout.usesSplitView + ? () => splitLeftHeaderItems + : canGoBack + ? undefined + : () => compactHomeHeaderItems + : undefined, + // Search lives in the persistent sidebar, so the split header keeps + // the git controls on the RIGHT (no center items — center space is + // reserved for future breadcrumbs/status). + unstable_headerRightItems: + Platform.OS === "ios" + ? () => (layout.usesSplitView ? threadCenterHeaderItems : compactRightHeaderItems) + : undefined, + unstable_headerSubtitle: usesNativeHeaderGlass ? headerSubtitle : undefined, + }), + [ + canGoBack, + compactHomeHeaderItems, + compactRightHeaderItems, + headerSubtitle, + layout.usesSplitView, + selectedThreadTitle, + splitLeftHeaderItems, + threadCenterHeaderItems, + usesNativeHeaderGlass, + ], + ); + if (!environmentId || !threadId) { return ; } @@ -788,8 +840,6 @@ function ThreadRouteContent( activePendingUserInputDrafts={requests.activePendingUserInputDrafts} activePendingUserInputAnswers={requests.activePendingUserInputAnswers} respondingUserInputId={requests.respondingUserInputId} - draftMessage={composer.draftMessage} - draftAttachments={composer.draftAttachments} connectionStateLabel={routeConnectionState} threadSyncStatus={selectedThreadDetailState.status} loadEarlier={loadEarlierTurns} @@ -832,41 +882,7 @@ function ThreadRouteContent( return ( <> {activeInspectorRenderer ? : null} - splitLeftHeaderItems - : canGoBack - ? undefined - : () => compactHomeHeaderItems - : undefined, - // Search lives in the persistent sidebar, so the split header keeps - // the git controls on the RIGHT (no center items — center space is - // reserved for future breadcrumbs/status). - unstable_headerRightItems: - Platform.OS === "ios" - ? () => (layout.usesSplitView ? threadCenterHeaderItems : compactRightHeaderItems) - : undefined, - unstable_headerSubtitle: usesNativeHeaderGlass ? headerSubtitle : undefined, - }} - /> + {Platform.OS === "android" ? ( onDeleteThread(thread), [onDeleteThread, thread]); const handleRegenerateTitle = useCallback( @@ -888,7 +893,7 @@ export const ThreadListV2Row = memo(function ThreadListV2Row(props: { > {snoozedRow && props.snoozeWakeLabelText !== undefined ? props.snoozeWakeLabelText - : relativeTime(thread.latestUserMessageAt ?? thread.updatedAt ?? thread.createdAt)} + : timeLabel} diff --git a/apps/mobile/src/features/threads/thread-work-log.tsx b/apps/mobile/src/features/threads/thread-work-log.tsx index a5adacb8d19b..a7ec703ff788 100644 --- a/apps/mobile/src/features/threads/thread-work-log.tsx +++ b/apps/mobile/src/features/threads/thread-work-log.tsx @@ -1,4 +1,5 @@ import * as Haptics from "expo-haptics"; +import { memo } from "react"; import { type AppSymbolName, SymbolView } from "../../components/AppSymbol"; import { LayoutAnimation, Pressable, ScrollView, View } from "react-native"; @@ -120,14 +121,40 @@ export function collapsedWorkLogHeight( ); } -export function ThreadWorkLog(props: { +type ThreadWorkLogProps = { readonly activities: ReadonlyArray; readonly copiedRowId: string | null; readonly expandedRows: Readonly>; readonly iconSubtleColor: import("react-native").ColorValue; readonly onCopyRow: (rowId: string, value: string) => void; readonly onToggleRow: (rowId: string) => void; -}) { +}; + +/** + * Copy feedback and row expansion are feed-wide maps, so a work log only + * re-renders when a change touches one of its own rows. + */ +export function threadWorkLogPropsEqual( + previous: ThreadWorkLogProps, + next: ThreadWorkLogProps, +): boolean { + if ( + previous.activities !== next.activities || + previous.iconSubtleColor !== next.iconSubtleColor || + previous.onCopyRow !== next.onCopyRow || + previous.onToggleRow !== next.onToggleRow + ) { + return false; + } + for (const activity of next.activities) { + const id = activity.id; + if ((previous.copiedRowId === id) !== (next.copiedRowId === id)) return false; + if ((previous.expandedRows[id] ?? false) !== (next.expandedRows[id] ?? false)) return false; + } + return true; +} + +export const ThreadWorkLog = memo(function ThreadWorkLog(props: ThreadWorkLogProps) { const pressedBackground = useThemeColor("--color-subtle"); const rows = visibleWorkLogActivities(props.activities).map((activity) => ({ ...activity, @@ -272,7 +299,7 @@ export function ThreadWorkLog(props: { ); -} +}, threadWorkLogPropsEqual); export function ThreadWorkGroupToggle(props: { readonly expanded: boolean; diff --git a/apps/mobile/src/features/threads/threadContentPresentation.test.ts b/apps/mobile/src/features/threads/threadContentPresentation.test.ts index c3524a587197..a083535302bd 100644 --- a/apps/mobile/src/features/threads/threadContentPresentation.test.ts +++ b/apps/mobile/src/features/threads/threadContentPresentation.test.ts @@ -70,3 +70,16 @@ describe("thread content presentation", () => { }); }); }); + +describe("projectThreadContentPresentation identity", () => { + it("returns the same object for equal inputs so memoized consumers skip renders", () => { + const ready = { hasDetail: true, detailError: null, detailDeleted: false } as const; + expect(projectThreadContentPresentation({ ...ready, connectionState: "connected" })).toBe( + projectThreadContentPresentation({ ...ready, connectionState: "connected" }), + ); + const failed = { hasDetail: false, detailError: "boom", detailDeleted: false } as const; + expect(projectThreadContentPresentation({ ...failed, connectionState: "connected" })).toBe( + projectThreadContentPresentation({ ...failed, connectionState: "connected" }), + ); + }); +}); diff --git a/apps/mobile/src/features/threads/threadContentPresentation.ts b/apps/mobile/src/features/threads/threadContentPresentation.ts index c755e4c79ff0..07031bb2337a 100644 --- a/apps/mobile/src/features/threads/threadContentPresentation.ts +++ b/apps/mobile/src/features/threads/threadContentPresentation.ts @@ -9,6 +9,22 @@ export type ThreadContentPresentation = readonly detail: string; }; +// Shared instances keep the presentation referentially stable across renders, +// so memoized consumers such as the thread feed skip unrelated re-renders. +const READY: ThreadContentPresentation = { kind: "ready" }; +const LOADING: ThreadContentPresentation = { kind: "loading" }; +const DELETED: ThreadContentPresentation = { + kind: "unavailable", + title: "Chat unavailable", + detail: "This chat was deleted or is no longer available.", +}; +const NOT_CACHED: ThreadContentPresentation = { + kind: "unavailable", + title: "Messages not cached", + detail: "Reconnect this environment to load the conversation.", +}; +const detailErrorPresentations = new Map(); + export function projectThreadContentPresentation(input: { readonly hasDetail: boolean; readonly detailError: string | null; @@ -16,21 +32,24 @@ export function projectThreadContentPresentation(input: { readonly connectionState: EnvironmentConnectionPhase; }): ThreadContentPresentation { if (input.hasDetail) { - return { kind: "ready" }; + return READY; } if (input.detailDeleted) { - return { - kind: "unavailable", - title: "Chat unavailable", - detail: "This chat was deleted or is no longer available.", - }; + return DELETED; } if (input.detailError !== null) { - return { - kind: "unavailable", - title: "Could not load conversation", - detail: input.detailError, - }; + let presentation = detailErrorPresentations.get(input.detailError); + if (presentation === undefined) { + presentation = { + kind: "unavailable", + title: "Could not load conversation", + detail: input.detailError, + }; + // Error text is unbounded; keep only the latest to avoid a slow leak. + detailErrorPresentations.clear(); + detailErrorPresentations.set(input.detailError, presentation); + } + return presentation; } if ( input.connectionState === "connected" || @@ -39,11 +58,7 @@ export function projectThreadContentPresentation(input: { ) { // Messages will arrive once the (re)connection completes — present as // loading; the composer's connection pill reports the connection phase. - return { kind: "loading" }; + return LOADING; } - return { - kind: "unavailable", - title: "Messages not cached", - detail: "Reconnect this environment to load the conversation.", - }; + return NOT_CACHED; } diff --git a/apps/mobile/src/lib/threadActivity.test.ts b/apps/mobile/src/lib/threadActivity.test.ts index 56ea095476f0..c9d113c815f7 100644 --- a/apps/mobile/src/lib/threadActivity.test.ts +++ b/apps/mobile/src/lib/threadActivity.test.ts @@ -19,7 +19,9 @@ import { deriveThreadFeedPresentation, isPendingUserInputOptionSelected, setPendingUserInputCustomAnswer, + threadFeedEntriesEqual, togglePendingUserInputOptionSelection, + unchangedPrefixLength, type ThreadFeedActivity, type ThreadFeedEntry, } from "./threadActivity"; @@ -295,6 +297,85 @@ describe("buildThreadFeed", () => { ]); }); + it("keeps unchanged rows referentially stable while a turn streams", () => { + const turnId = TurnId.make("turn-stream"); + const user = { + id: MessageId.make("user-stream"), + role: "user" as const, + text: "Run the tests", + turnId: null, + streaming: false, + createdAt: "2026-08-31T00:00:00.000Z", + updatedAt: "2026-08-31T00:00:00.000Z", + }; + const toolUpdated = makeActivity({ + id: EventId.make("stream-tool-updated"), + kind: "tool.updated", + tone: "tool", + summary: "Run tests", + createdAt: "2026-08-31T00:00:01.000Z", + turnId, + payload: { title: "Run tests", itemType: "command_execution", detail: "bun run test" }, + }); + const assistant = { + id: MessageId.make("assistant-stream"), + role: "assistant" as const, + text: "Work", + turnId, + streaming: true, + createdAt: "2026-08-31T00:00:02.000Z", + updatedAt: "2026-08-31T00:00:02.000Z", + }; + const base = { + id: ThreadId.make("thread-stream"), + projectId: ProjectId.make("project-1"), + title: "Streaming", + }; + const first = buildThreadFeed( + makeThread({ ...base, messages: [user, assistant], activities: [toolUpdated] }), + ); + const grown = { ...assistant, text: "Working", updatedAt: "2026-08-31T00:00:03.000Z" }; + const second = buildThreadFeed( + makeThread({ ...base, messages: [user, grown], activities: [toolUpdated] }), + ); + + expect(second.map((entry) => entry.id)).toEqual(first.map((entry) => entry.id)); + expect(second[0]).toBe(first[0]); + const firstGroup = first[1]; + const secondGroup = second[1]; + expect(firstGroup?.type).toBe("activity-group"); + if (firstGroup?.type !== "activity-group" || secondGroup?.type !== "activity-group") return; + expect(secondGroup.activities[0]).toBe(firstGroup.activities[0]); + expect(second[2]).not.toBe(first[2]); + expect(second[2]).toMatchObject({ type: "message", message: { text: "Working" } }); + + // Presentation rebuilds wrapper rows; list equality still sees only the streamed row change. + const presentFirst = deriveThreadFeedPresentation(first, null, new Set()); + const presentSecond = deriveThreadFeedPresentation(second, null, new Set()); + expect(presentSecond[1]).not.toBe(presentFirst[1]); + expect( + presentSecond.map((entry, index) => threadFeedEntriesEqual(presentFirst[index]!, entry)), + ).toEqual([true, true, false]); + + // A lifecycle row merged from a new activity gets a fresh presentation. + const toolCompleted = makeActivity({ + ...toolUpdated, + id: EventId.make("stream-tool-completed"), + kind: "tool.completed", + summary: "Run tests completed", + createdAt: "2026-08-31T00:00:01.500Z", + }); + const third = buildThreadFeed( + makeThread({ ...base, messages: [user, grown], activities: [toolUpdated, toolCompleted] }), + ); + const thirdGroup = third[1]; + if (thirdGroup?.type !== "activity-group") throw new Error("expected activity group"); + expect(thirdGroup.activities).toHaveLength(1); + expect(thirdGroup.activities[0]).not.toBe(firstGroup.activities[0]); + expect(thirdGroup.activities[0]?.id).toBe("stream-tool-completed"); + expect(threadFeedEntriesEqual(secondGroup, thirdGroup)).toBe(false); + }); + it("keeps older local feedback before newer messages returned by the server", () => { const submission = { id: MessageId.make("feedback-command-ordering"), @@ -904,3 +985,25 @@ describe("quiet timeline: nested agents", () => { expect(ids).not.toContain("shell-done"); }); }); + +describe("unchangedPrefixLength", () => { + it("reports the whole previous length for appends", () => { + const a = { id: "a" }; + const b = { id: "b" }; + const c = { id: "c" }; + expect(unchangedPrefixLength([], [a])).toBe(0); + expect(unchangedPrefixLength([a, b], [a, b, c])).toBe(2); + expect(unchangedPrefixLength([a, b], [a, b])).toBe(2); + }); + + it("stops at the first removed, inserted, or replaced item", () => { + const a = { id: "a" }; + const b = { id: "b" }; + const c = { id: "c" }; + expect(unchangedPrefixLength([a, b, c], [a, c])).toBe(1); + expect(unchangedPrefixLength([a, b], [a, c, b])).toBe(1); + expect(unchangedPrefixLength([a, b], [{ id: "a" }, b])).toBe(0); + // Interior replacement by id keeps both endpoints but must still rescan it. + expect(unchangedPrefixLength([a, b, c], [a, { id: "b" }, c])).toBe(1); + }); +}); diff --git a/apps/mobile/src/lib/threadActivity.ts b/apps/mobile/src/lib/threadActivity.ts index b09a1a6c34df..f388a728390c 100644 --- a/apps/mobile/src/lib/threadActivity.ts +++ b/apps/mobile/src/lib/threadActivity.ts @@ -232,7 +232,7 @@ function isAgentInternalActivity(activity: OrchestrationThreadActivity): boolean function deriveWorkLogEntries( activities: ReadonlyArray, -): DerivedWorkLogEntry[] { +): CollapsedWorkLogEntry[] { const ordered = Arr.sort(activities, activityOrder); const entries: DerivedWorkLogEntry[] = []; for (const activity of ordered) { @@ -247,11 +247,24 @@ function deriveWorkLogEntries( if (activity.summary === "Checkpoint captured") continue; if (isPlanBoundaryToolActivity(activity)) continue; if (isAgentInternalActivity(activity)) continue; - entries.push(toDerivedWorkLogEntry(activity)); + entries.push(cachedDerivedWorkLogEntry(activity)); } return collapseDerivedWorkLogEntries(entries); } +// Activities are immutable and keep their identity across stream updates, so +// each one is parsed into a work-log entry once instead of on every delta. +const derivedWorkLogEntryCache = new WeakMap(); + +function cachedDerivedWorkLogEntry(activity: OrchestrationThreadActivity): DerivedWorkLogEntry { + let entry = derivedWorkLogEntryCache.get(activity); + if (entry === undefined) { + entry = toDerivedWorkLogEntry(activity); + derivedWorkLogEntryCache.set(activity, entry); + } + return entry; +} + function isPlanBoundaryToolActivity(activity: OrchestrationThreadActivity): boolean { if (activity.kind !== "tool.updated" && activity.kind !== "tool.completed") { return false; @@ -360,10 +373,20 @@ function toDerivedWorkLogEntry(activity: OrchestrationThreadActivity): DerivedWo return entry; } +/** + * A collapsed row plus the per-activity entries it was merged from. The sources + * are identity-stable, so they tell whether a row's inputs changed. + */ +interface CollapsedWorkLogEntry { + readonly entry: DerivedWorkLogEntry; + readonly sources: ReadonlyArray; +} + function collapseDerivedWorkLogEntries( entries: ReadonlyArray, -): DerivedWorkLogEntry[] { +): CollapsedWorkLogEntry[] { const collapsed: DerivedWorkLogEntry[] = []; + const sources: DerivedWorkLogEntry[][] = []; // Subagent rows collapse by identity, not adjacency (quiet-timeline // guarantee; mirrors web's session-logic). const taskRowIndex = new Map(); @@ -377,20 +400,24 @@ function collapseDerivedWorkLogEntries( const existingIndex = taskRowIndex.get(entry.taskId); if (existingIndex !== undefined) { collapsed[existingIndex] = mergeDerivedWorkLogEntries(collapsed[existingIndex]!, entry); + sources[existingIndex]!.push(entry); continue; } taskRowIndex.set(entry.taskId, collapsed.length); collapsed.push(entry); + sources.push([entry]); continue; } const previous = collapsed.at(-1); if (previous && shouldCollapseToolLifecycleEntries(previous, entry)) { collapsed[collapsed.length - 1] = mergeDerivedWorkLogEntries(previous, entry); + sources[sources.length - 1]!.push(entry); continue; } collapsed.push(entry); + sources.push([entry]); } - return collapsed; + return collapsed.map((entry, index) => ({ entry, sources: sources[index]! })); } function shouldCollapseToolLifecycleEntries( @@ -1368,59 +1395,203 @@ export function buildThreadFeed( options?.loadedMessages !== undefined ? (loadedMessages[0]?.createdAt ?? null) : null; const botStepMeters = buildBotStepMeters(thread.activities); const workLogEntries = deriveWorkLogEntries(thread.activities); - const entries = Arr.sortWith( - [ - ...messages.map((message) => ({ - type: "message", - id: message.id, - createdAt: message.createdAt, - message, - ...(message.turnId === null ? {} : { botStepMeter: botStepMeters.get(message.turnId) }), - })), - ...workLogEntries - .filter((entry) => { - if (options?.loadedMessages === undefined) { - return true; - } - return ( - oldestLoadedMessageCreatedAt === null || entry.createdAt >= oldestLoadedMessageCreatedAt - ); - }) - .map((entry) => { - const summary = workEntryHeading(entry); - const detail = workEntryPreview(entry); - const getFullDetail = memoizeValue(() => buildWorkEntryExpandedBody(entry)); - const getCopyText = memoizeValue(() => - [summary, detail, getFullDetail()] - .filter((value, index, values): value is string => { - return Boolean(value) && values.indexOf(value) === index; - }) - .join("\n"), - ); - return { - type: "activity", - id: entry.id, - createdAt: entry.createdAt, - turnId: entry.turnId, - activity: { - id: entry.id, - createdAt: entry.createdAt, - turnId: entry.turnId, - summary, - detail, - canExpand: workEntryHasExpandedBody(entry), - getFullDetail, - getCopyText, - icon: workEntryIcon(entry), - toolLike: workLogEntryIsToolLike(entry), - status: workEntryStatus(entry), - }, - }; - }), - ], - (s) => new Date(s.createdAt), - Order.Date, + const timed: Array<{ readonly at: number; readonly entry: RawThreadFeedEntry }> = []; + for (const message of messages) { + const entry = messageFeedEntry( + message, + message.turnId === null ? undefined : botStepMeters.get(message.turnId), + ); + timed.push({ at: Date.parse(entry.createdAt), entry }); + } + for (const collapsed of workLogEntries) { + if ( + options?.loadedMessages !== undefined && + oldestLoadedMessageCreatedAt !== null && + collapsed.entry.createdAt < oldestLoadedMessageCreatedAt + ) { + continue; + } + const entry = activityFeedEntry(collapsed); + timed.push({ at: Date.parse(entry.createdAt), entry }); + } + // Timestamps are parsed once above; sorting on numbers avoids allocating a + // Date per comparison. + timed.sort((left, right) => compareFeedTimes(left.at, right.at)); + + return groupAdjacentActivities(timed.map((item) => item.entry)); +} + +/** Same ordering as `Order.Date`: stable ties, unparsable times first. */ +function compareFeedTimes(left: number, right: number): number { + if (left === right) return 0; + const leftInvalid = Number.isNaN(left); + const rightInvalid = Number.isNaN(right); + if (leftInvalid && rightInvalid) return 0; + if (leftInvalid) return -1; + if (rightInvalid) return 1; + return left < right ? -1 : 1; +} + +type MessageFeedEntry = Extract; + +// Feed entries are cached on their source objects so unchanged rows keep their +// identity while a turn streams, letting the list skip re-rendering them. +const messageFeedEntryCache = new WeakMap< + OrchestrationThread["messages"][number], + MessageFeedEntry +>(); + +function messageFeedEntry( + message: OrchestrationThread["messages"][number], + botStepMeter: BotStepMeterData | undefined, +): MessageFeedEntry { + const cached = messageFeedEntryCache.get(message); + if (cached !== undefined && botStepMetersEqual(cached.botStepMeter, botStepMeter)) { + return cached; + } + const entry: MessageFeedEntry = { + type: "message", + id: message.id, + createdAt: message.createdAt, + message, + ...(message.turnId === null ? {} : { botStepMeter }), + }; + messageFeedEntryCache.set(message, entry); + return entry; +} + +function botStepMetersEqual( + left: BotStepMeterData | undefined, + right: BotStepMeterData | undefined, +): boolean { + if (left === right) return true; + if (left === undefined || right === undefined) return false; + return ( + left.tokens === right.tokens && + left.costUsd === right.costUsd && + left.hardStopReached === right.hardStopReached && + (left.engine === right.engine || + (left.engine.provider === right.engine.provider && + left.engine.model === right.engine.model && + left.engine.options === right.engine.options)) + ); +} + +type ActivityFeedEntry = Extract; + +// Keyed on the newest source entry; reused only when every source matches. +const activityFeedEntryCache = new WeakMap< + DerivedWorkLogEntry, + { readonly sources: ReadonlyArray; readonly entry: ActivityFeedEntry } +>(); + +function activityFeedEntry(collapsed: CollapsedWorkLogEntry): ActivityFeedEntry { + const key = collapsed.sources[collapsed.sources.length - 1]!; + const cached = activityFeedEntryCache.get(key); + if (cached !== undefined && sameEntries(cached.sources, collapsed.sources)) { + return cached.entry; + } + const entry = toActivityFeedEntry(collapsed.entry); + activityFeedEntryCache.set(key, { sources: collapsed.sources, entry }); + return entry; +} + +function sameEntries(left: ReadonlyArray, right: ReadonlyArray): boolean { + if (left.length !== right.length) return false; + for (let index = 0; index < left.length; index += 1) { + if (left[index] !== right[index]) return false; + } + return true; +} + +function toActivityFeedEntry(entry: DerivedWorkLogEntry): ActivityFeedEntry { + const summary = workEntryHeading(entry); + const detail = workEntryPreview(entry); + const getFullDetail = memoizeValue(() => buildWorkEntryExpandedBody(entry)); + const getCopyText = memoizeValue(() => + [summary, detail, getFullDetail()] + .filter((value, index, values): value is string => { + return Boolean(value) && values.indexOf(value) === index; + }) + .join("\n"), ); + return { + type: "activity", + id: entry.id, + createdAt: entry.createdAt, + turnId: entry.turnId, + activity: { + id: entry.id, + createdAt: entry.createdAt, + turnId: entry.turnId, + summary, + detail, + canExpand: workEntryHasExpandedBody(entry), + getFullDetail, + getCopyText, + icon: workEntryIcon(entry), + toolLike: workLogEntryIsToolLike(entry), + status: workEntryStatus(entry), + }, + }; +} + +/** + * Length of the leading run of `previous` that `next` still holds by identity. + * Appends keep the whole previous length; a removal or replacement stops at the + * first changed item so the caller rescans from there. + */ +export function unchangedPrefixLength( + previous: ReadonlyArray, + next: ReadonlyArray, +): number { + const length = Math.min(previous.length, next.length); + let index = 0; + while (index < length && previous[index] === next[index]) { + index += 1; + } + return index; +} - return groupAdjacentActivities(entries); +/** + * Row equality for the thread feed list. Presentation rebuilds wrapper rows on + * every feed change, so the list compares their rendered fields instead of + * identity and keeps unchanged rows mounted without re-rendering them. + */ +export function threadFeedEntriesEqual(previous: ThreadFeedEntry, next: ThreadFeedEntry): boolean { + if (previous === next) return true; + if (previous.id !== next.id || previous.createdAt !== next.createdAt) return false; + switch (previous.type) { + case "message": + return ( + next.type === "message" && + previous.message === next.message && + botStepMetersEqual(previous.botStepMeter, next.botStepMeter) + ); + case "working": + return next.type === "working"; + case "activity-group": + return ( + next.type === "activity-group" && + previous.turnId === next.turnId && + previous.activities.length === next.activities.length && + previous.activities.every((activity, index) => activity === next.activities[index]) + ); + case "work-toggle": + return ( + next.type === "work-toggle" && + previous.turnId === next.turnId && + previous.groupId === next.groupId && + previous.hiddenCount === next.hiddenCount && + previous.expanded === next.expanded && + previous.onlyToolActivities === next.onlyToolActivities + ); + case "turn-fold": + return ( + next.type === "turn-fold" && + previous.turnId === next.turnId && + previous.label === next.label && + previous.expanded === next.expanded + ); + } } diff --git a/apps/mobile/src/native/StackHeader.tsx b/apps/mobile/src/native/StackHeader.tsx index a524d5152473..536e271a0b05 100644 --- a/apps/mobile/src/native/StackHeader.tsx +++ b/apps/mobile/src/native/StackHeader.tsx @@ -157,14 +157,20 @@ export function NativeStackScreenOptions(props: { const latestOptionFunctionsRef = useRef(new Map unknown>()); const optionFunctionWrappersRef = useRef(new Map unknown>()); const normalizedOptions = useMemo(() => normalizeScreenOptions(props.options), [props.options]); - const stableOptions = normalizedOptions - ? (stabilizeOptionFunctions( - normalizedOptions, - "options", - latestOptionFunctionsRef.current, - optionFunctionWrappersRef.current, - ) as NativeStackNavigationOptions) - : undefined; + // Keyed on the options identity: callers that memoize their options skip the + // deep copy and the signature walk below on unrelated re-renders. + const stableOptions = useMemo( + () => + normalizedOptions + ? (stabilizeOptionFunctions( + normalizedOptions, + "options", + latestOptionFunctionsRef.current, + optionFunctionWrappersRef.current, + ) as NativeStackNavigationOptions) + : undefined, + [normalizedOptions], + ); useLayoutEffect(() => { if (!navigation || !stableOptions) { diff --git a/apps/mobile/src/state/use-composer-drafts.ts b/apps/mobile/src/state/use-composer-drafts.ts index 3991cc8733ef..aeb9957fcc78 100644 --- a/apps/mobile/src/state/use-composer-drafts.ts +++ b/apps/mobile/src/state/use-composer-drafts.ts @@ -10,7 +10,7 @@ import { type RuntimeMode, } from "@t3tools/contracts"; import * as Schema from "effect/Schema"; -import { useEffect } from "react"; +import { useEffect, useMemo } from "react"; import { Atom } from "effect/unstable/reactivity"; import { writeFileAtomically } from "../lib/atomic-file"; @@ -102,6 +102,25 @@ export const composerDraftsAtom = Atom.make>({}).p Atom.withLabel("mobile:composer-drafts"), ); +// Per-key view of the draft map. The derived atom only notifies when this +// key's draft object changes, so typing in one chat does not re-render +// components that read another chat's draft. +const composerDraftAtom = Atom.family((draftKey: string) => + Atom.make((get): ComposerDraft | undefined => get(composerDraftsAtom)[draftKey]), +); + +function composerDraftFieldAtomFamily(field: K) { + return Atom.family((draftKey: string) => + Atom.make((get): ComposerDraft[K] | undefined => get(composerDraftsAtom)[draftKey]?.[field]), + ); +} + +// Draft edits spread the previous draft, so these references only change when +// the setting itself changes, not on every keystroke. +const composerDraftModelSelectionAtom = composerDraftFieldAtomFamily("modelSelection"); +const composerDraftRuntimeModeAtom = composerDraftFieldAtomFamily("runtimeMode"); +const composerDraftInteractionModeAtom = composerDraftFieldAtomFamily("interactionMode"); + let loadPromise: Promise | null = null; let persistTimer: ReturnType | null = null; let persistRetryNeeded = false; @@ -685,9 +704,22 @@ export async function clearComposerDraftsEnvironment(environmentId: EnvironmentI } export function useComposerDraft(draftKey: string | null): ComposerDraft { - const drafts = useAtomValue(composerDraftsAtom); + const draft = useAtomValue(composerDraftAtom(draftKey ?? "")); useEffect(() => { ensureComposerDraftsLoaded(); }, []); - return draftKey ? normalizeDraft(drafts[draftKey]) : EMPTY_DRAFT; + return useMemo(() => (draftKey ? normalizeDraft(draft) : EMPTY_DRAFT), [draft, draftKey]); +} + +/** Reads a draft's settings without subscribing to its text or attachments. */ +export function useComposerDraftSettings(draftKey: string | null): { + readonly modelSelection: ModelSelection | undefined; + readonly runtimeMode: RuntimeMode | undefined; + readonly interactionMode: ProviderInteractionMode | undefined; +} { + const key = draftKey ?? ""; + const modelSelection = useAtomValue(composerDraftModelSelectionAtom(key)); + const runtimeMode = useAtomValue(composerDraftRuntimeModeAtom(key)); + const interactionMode = useAtomValue(composerDraftInteractionModeAtom(key)); + return { modelSelection, runtimeMode, interactionMode }; } diff --git a/apps/mobile/src/state/use-thread-composer-state.ts b/apps/mobile/src/state/use-thread-composer-state.ts index 6b39af6137aa..90593f18393b 100644 --- a/apps/mobile/src/state/use-thread-composer-state.ts +++ b/apps/mobile/src/state/use-thread-composer-state.ts @@ -1,4 +1,3 @@ -import { useAtomValue } from "@effect/atom-react"; import { useCallback, useEffect, useMemo, useRef, useState } from "react"; import { Alert } from "react-native"; import * as Cause from "effect/Cause"; @@ -31,7 +30,7 @@ import { import type { DraftComposerImageAttachment } from "../lib/composerImages"; import { scopedThreadKey } from "../lib/scopedEntities"; import { copyTextWithHaptic } from "../lib/copyTextWithHaptic"; -import { buildThreadFeed } from "../lib/threadActivity"; +import { buildThreadFeed, unchangedPrefixLength } from "../lib/threadActivity"; import { tryOpenExternalUrl } from "../lib/openExternalUrl"; import { appAtomRegistry } from "../state/atom-registry"; import { @@ -46,6 +45,7 @@ import { setComposerDraftText, updateComposerDraftSettings, useComposerDraft, + useComposerDraftSettings, } from "./use-composer-drafts"; import { setPendingConnectionError } from "../state/use-remote-environment-registry"; import { useSelectedThreadDetail } from "../state/use-thread-detail"; @@ -90,7 +90,6 @@ export function useThreadComposerState() { const { selectedThread: selectedThreadShell, selectedEnvironmentRuntime } = useThreadSelection(); const selectedThreadDetail = useSelectedThreadDetail(); const openedAuthorizationActivitiesRef = useRef(new Set()); - const composerDrafts = useAtomValue(composerDraftsAtom); const queuedMessagesByThreadKey = useThreadOutboxMessages(); const [feedbackSubmissionsByThreadKey, setFeedbackSubmissionsByThreadKey] = useState< Record> @@ -103,9 +102,15 @@ export function useThreadComposerState() { ensureComposerDraftsLoaded(); }, []); + const scannedAuthorizationActivitiesRef = useRef>([]); useEffect(() => { - for (const activity of selectedThreadDetail?.activities ?? []) { - if (activity.kind !== "mcp.oauth.authorization-required") continue; + const activities = selectedThreadDetail?.activities ?? []; + // Activity updates usually append, so only the unseen tail is scanned. + const start = unchangedPrefixLength(scannedAuthorizationActivitiesRef.current, activities); + scannedAuthorizationActivitiesRef.current = activities; + for (let index = start; index < activities.length; index += 1) { + const activity = activities[index]; + if (activity?.kind !== "mcp.oauth.authorization-required") continue; if (openedAuthorizationActivitiesRef.current.has(activity.id)) continue; if (!activity.payload || typeof activity.payload !== "object") continue; const authorizationUrl = (activity.payload as Record).authorizationUrl; @@ -117,6 +122,8 @@ export function useThreadComposerState() { } openedAuthorizationActivitiesRef.current.add(activity.id); void tryOpenExternalUrl(authorizationUrl, "mcp-oauth"); + // Leave later rows unscanned so the next update can open them. + scannedAuthorizationActivitiesRef.current = activities.slice(0, index + 1); break; } }, [selectedThreadDetail?.activities]); @@ -144,14 +151,14 @@ export function useThreadComposerState() { }); }, [feedbackSubmissionsByThreadKey, selectedThreadDetail, selectedThreadKey]); - const selectedDraft = selectedThreadKey ? composerDrafts[selectedThreadKey] : null; - const draftMessage = selectedDraft?.text ?? ""; - const draftAttachments = selectedDraft?.attachments ?? []; + // Draft text and attachments are read by the composer itself, so typing + // re-renders the composer rather than the whole thread route. + const selectedDraft = useComposerDraftSettings(selectedThreadKey); const selectedThreadQueueCount = selectedThreadQueuedMessages.length; const selectedThread = selectedThreadDetail ?? selectedThreadShell; - const modelSelection = selectedDraft?.modelSelection ?? selectedThread?.modelSelection ?? null; - const runtimeMode = selectedDraft?.runtimeMode ?? selectedThread?.runtimeMode ?? null; - const interactionMode = selectedDraft?.interactionMode ?? selectedThread?.interactionMode ?? null; + const modelSelection = selectedDraft.modelSelection ?? selectedThread?.modelSelection ?? null; + const runtimeMode = selectedDraft.runtimeMode ?? selectedThread?.runtimeMode ?? null; + const interactionMode = selectedDraft.interactionMode ?? selectedThread?.interactionMode ?? null; const selectedThreadSessionActivity = useMemo(() => { const selectedThread = selectedThreadDetail ?? selectedThreadShell; @@ -314,7 +321,7 @@ export function useThreadComposerState() { const threadKey = scopedThreadKey(selectedThreadShell.environmentId, selectedThreadShell.id); const result = await pickComposerImages({ - existingCount: composerDrafts[threadKey]?.attachments.length ?? 0, + existingCount: getComposerDraftSnapshot(threadKey).attachments.length, }); if (result.images.length > 0) { appendComposerDraftAttachments(threadKey, result.images); @@ -322,7 +329,7 @@ export function useThreadComposerState() { if (result.error) { setPendingConnectionError(result.error); } - }, [composerDrafts, selectedThreadShell]); + }, [selectedThreadShell]); const onPasteIntoDraft = useCallback(async () => { if (!selectedThreadShell) { @@ -331,7 +338,7 @@ export function useThreadComposerState() { const threadKey = scopedThreadKey(selectedThreadShell.environmentId, selectedThreadShell.id); const result = await pasteComposerClipboard({ - existingCount: composerDrafts[threadKey]?.attachments.length ?? 0, + existingCount: getComposerDraftSnapshot(threadKey).attachments.length, }); if (result.images.length > 0) { appendComposerDraftAttachments(threadKey, result.images); @@ -342,7 +349,7 @@ export function useThreadComposerState() { if (result.error) { setPendingConnectionError(result.error); } - }, [composerDrafts, selectedThreadShell]); + }, [selectedThreadShell]); const onNativePasteImages = useCallback( async (uris: ReadonlyArray) => { @@ -354,7 +361,7 @@ export function useThreadComposerState() { try { const images = await convertPastedImagesToAttachments({ uris, - existingCount: composerDrafts[threadKey]?.attachments.length ?? 0, + existingCount: getComposerDraftSnapshot(threadKey).attachments.length, }); if (images.length > 0) { appendComposerDraftAttachments(threadKey, images); @@ -368,7 +375,7 @@ export function useThreadComposerState() { }); } }, - [composerDrafts, selectedThreadShell], + [selectedThreadShell], ); const onRemoveDraftImage = useCallback( @@ -417,8 +424,6 @@ export function useThreadComposerState() { selectedThreadFeed, selectedThreadQueueCount, activeWorkStartedAt, - draftMessage, - draftAttachments, modelSelection, runtimeMode, interactionMode, diff --git a/packages/client-runtime/src/state/threadReducer.test.ts b/packages/client-runtime/src/state/threadReducer.test.ts index f2abec2e8a5f..e36e5f845fb4 100644 --- a/packages/client-runtime/src/state/threadReducer.test.ts +++ b/packages/client-runtime/src/state/threadReducer.test.ts @@ -1190,6 +1190,82 @@ describe("applyThreadDetailEvent", () => { } }); + it("supersedes context-window updates on the in-order append path", () => { + const activity = ( + id: string, + sequence: number, + kind: string, + turn: string, + usedTokens?: number, + ) => ({ + id: EventId.make(id), + tone: "info" as const, + kind, + summary: id, + payload: usedTokens === undefined ? {} : { usedTokens }, + turnId: TurnId.make(turn), + sequence, + createdAt: "2026-04-01T11:00:00.000Z", + }); + const append = ( + thread: OrchestrationThread, + sequence: number, + next: ReturnType, + ) => { + const result = applyThreadDetailEvent(thread, { + ...baseEventFields, + sequence, + occurredAt: "2026-04-01T11:02:00.000Z", + aggregateKind: "thread", + aggregateId: ThreadId.make("thread-1"), + type: "thread.activity-appended", + payload: { threadId: ThreadId.make("thread-1"), activity: next }, + }); + assert(result.kind === "updated"); + return result.thread; + }; + + // The first append sorts the snapshot array and indexes it; the later + // ones take the in-order path. + let thread = append( + { + ...baseThread, + activities: [ + activity("cw-turn-0", 1, "context-window.updated", "turn-0", 100), + activity("cw-1", 2, "context-window.updated", "turn-1", 1_000), + ], + }, + 30, + activity("tool-1", 3, "command", "turn-1"), + ); + const beforeSupersede = thread; + thread = append(thread, 31, activity("cw-2", 4, "context-window.updated", "turn-1", 2_000)); + thread = append(thread, 32, activity("tool-2", 5, "command", "turn-1")); + thread = append(thread, 33, activity("cw-3", 6, "context-window.updated", "turn-1", 3_000)); + + expect(thread.activities.map((entry) => entry.id)).toEqual([ + "cw-turn-0", + "tool-1", + "tool-2", + "cw-3", + ]); + // The input thread is not mutated. + expect(beforeSupersede.activities.map((entry) => entry.id)).toEqual([ + "cw-turn-0", + "cw-1", + "tool-1", + ]); + // A re-delivered current row still dedupes rather than duplicating. + thread = append(thread, 34, activity("cw-3", 7, "context-window.updated", "turn-1", 3_500)); + expect(thread.activities.map((entry) => entry.id)).toEqual([ + "cw-turn-0", + "tool-1", + "tool-2", + "cw-3", + ]); + expect(thread.activities.at(-1)?.payload).toEqual({ usedTokens: 3_500 }); + }); + it("does not collapse context-window history for a malformed update", () => { const resolvable = { id: EventId.make("activity-cw-resolvable"), diff --git a/packages/client-runtime/src/state/threadReducer.ts b/packages/client-runtime/src/state/threadReducer.ts index a1b3cbf94a84..44c11927c6de 100644 --- a/packages/client-runtime/src/state/threadReducer.ts +++ b/packages/client-runtime/src/state/threadReducer.ts @@ -650,12 +650,27 @@ export function applyThreadDetailEvent( const ids = activityIdIndex.get(thread.activities); const lastActivity = thread.activities.at(-1); if ( - !supersedesContextWindow && ids !== undefined && (lastActivity === undefined || activityOrder(lastActivity, activity) <= 0) && !ids.has(activity.id) ) { - const activities = Arr.append(thread.activities, activity); + let activities: ReadonlyArray; + if (supersedesContextWindow) { + // Dropping rows from a sorted array keeps it sorted, so a superseding + // update copies once and appends instead of re-sorting the history. + const retained: Array = []; + for (const entry of thread.activities) { + if (entry.turnId === activity.turnId && isResolvableContextWindowActivity(entry)) { + ids.delete(entry.id); + } else { + retained.push(entry); + } + } + retained.push(activity); + activities = retained; + } else { + activities = Arr.append(thread.activities, activity); + } activityIdIndex.delete(thread.activities); ids.add(activity.id); activityIdIndex.set(activities, ids); diff --git a/packages/client-runtime/src/state/threads-sync.test.ts b/packages/client-runtime/src/state/threads-sync.test.ts index 410de1015696..3debf79c02f6 100644 --- a/packages/client-runtime/src/state/threads-sync.test.ts +++ b/packages/client-runtime/src/state/threads-sync.test.ts @@ -99,7 +99,21 @@ const ACTIVE_THREAD: OrchestrationThread = { }, }; -type TestThreadInput = OrchestrationThreadStreamItem | Error; +// An array is delivered as one transport chunk, like a burst read off the socket. +type TestThreadInput = + | OrchestrationThreadStreamItem + | readonly [OrchestrationThreadStreamItem, ...OrchestrationThreadStreamItem[]] + | Error; + +type StreamBatch = readonly [OrchestrationThreadStreamItem, ...OrchestrationThreadStreamItem[]]; + +function toStreamBatch(input: Exclude): StreamBatch { + return isStreamBatch(input) ? input : [input]; +} + +function isStreamBatch(input: Exclude): input is StreamBatch { + return Array.isArray(input); +} function testSession( client: WsRpcProtocolClient, @@ -151,8 +165,9 @@ const makeHarness = Effect.fn("TestEnvironmentThreads.makeHarness")(function* (o const streamFrom = (queue: Queue.Queue) => Stream.fromQueue(queue).pipe( Stream.mapEffect((input) => - input instanceof Error ? Effect.fail(input) : Effect.succeed(input), + input instanceof Error ? Effect.fail(input) : Effect.succeed(toStreamBatch(input)), ), + Stream.flattenArray, ); const client = { [ORCHESTRATION_WS_METHODS.subscribeThread]: (input: { @@ -369,6 +384,36 @@ describe("EnvironmentThreads", () => { }), ); + it.effect("publishes a burst of events delivered together as one state change", () => + Effect.gen(function* () { + const harness = yield* makeHarness({ cached: BASE_THREAD }); + yield* awaitThreadState(harness.observed, (value) => value.status === "live"); + + yield* Queue.offer(harness.inputs, [ + titleUpdated("First", CACHED_SNAPSHOT_SEQUENCE + 1), + titleUpdated("Second", CACHED_SNAPSHOT_SEQUENCE + 2), + // A replayed sequence inside the batch is still ignored. + titleUpdated("Replayed", CACHED_SNAPSHOT_SEQUENCE + 2), + titleUpdated("Third", CACHED_SNAPSHOT_SEQUENCE + 3), + ]); + const published = yield* awaitThreadState( + harness.observed, + (value) => Option.isSome(value.data) && value.data.value.title !== BASE_THREAD.title, + ); + + // The first state carrying any event already carries all of them. + expect(Option.getOrThrow(published.data).title).toBe("Third"); + + // The resume cursor covers the whole batch. + yield* harness.replaceSession; + for (let attempt = 0; attempt < 100; attempt += 1) { + if ((yield* Ref.get(harness.subscriptionCount)) >= 2) break; + yield* Effect.yieldNow; + } + expect(yield* Ref.get(harness.lastSubscribeAfterSequence)).toBe(CACHED_SNAPSHOT_SEQUENCE + 3); + }), + ); + it.effect("does not persist active thread snapshots during streaming or teardown", () => Effect.gen(function* () { const savedThreads = yield* Effect.scoped( diff --git a/packages/client-runtime/src/state/threads.ts b/packages/client-runtime/src/state/threads.ts index fd97f693679d..b00375c1e837 100644 --- a/packages/client-runtime/src/state/threads.ts +++ b/packages/client-runtime/src/state/threads.ts @@ -313,62 +313,93 @@ export const makeEnvironmentThreadState = Effect.fn("EnvironmentThreadState.make ); }); - // Body of applyItem, running under applyLock. - const applyItemLocked = Effect.fn("EnvironmentThreadState.applyItemLocked")(function* ( - item: OrchestrationThreadStreamItem, + // Body of applyItems, running under applyLock. Event items reduce into a + // local thread that is published once per batch: a burst of streaming + // deltas delivered together costs one state change (and one UI render) + // instead of one per delta. The sequence cursor is committed together with + // the thread so a resubscribe never resumes past unpublished events. + const applyItemsLocked = Effect.fn("EnvironmentThreadState.applyItemsLocked")(function* ( + items: ReadonlyArray, ) { - if (item.kind === "synchronized") { - yield* Ref.set(awaitingCompletion, false); - yield* SubscriptionRef.update(state, (current) => - Option.isSome(current.data) && current.status !== "deleted" - ? { ...current, status: "live" as const, error: Option.none() } - : current, - ); - return; - } + let unpublishedThread: OrchestrationThread | null = null; + let unpublishedSequence: number | null = null; + const publish = Effect.gen(function* () { + if (unpublishedSequence !== null) { + yield* SubscriptionRef.set(lastSequence, unpublishedSequence); + unpublishedSequence = null; + } + if (unpublishedThread !== null) { + const thread = unpublishedThread; + unpublishedThread = null; + yield* setThread(thread, "keep"); + } + }); - if (item.kind === "snapshot") { - // A fresh snapshot replaces all loaded history, including older - // pages: a turn reverted while disconnected would otherwise survive - // in the preserved history with no event left to remove it. The - // epoch bump discards any older-page fetch racing this snapshot. - yield* Ref.update(historyEpoch, (epoch) => epoch + 1); - yield* SubscriptionRef.set(lastSequence, item.snapshot.snapshotSequence); - yield* setThread(item.snapshot.thread, pageStateFromSnapshot(item.snapshot.page)); - return; - } + for (const item of items) { + if (item.kind === "synchronized") { + yield* publish; + yield* Ref.set(awaitingCompletion, false); + yield* SubscriptionRef.update(state, (current) => + Option.isSome(current.data) && current.status !== "deleted" + ? { ...current, status: "live" as const, error: Option.none() } + : current, + ); + continue; + } - const sequence = yield* SubscriptionRef.get(lastSequence); - if (item.event.sequence <= sequence) { - return; - } - yield* SubscriptionRef.set(lastSequence, item.event.sequence); + if (item.kind === "snapshot") { + // A fresh snapshot replaces all loaded history, including older + // pages: a turn reverted while disconnected would otherwise survive + // in the preserved history with no event left to remove it. The + // epoch bump discards any older-page fetch racing this snapshot. + unpublishedThread = null; + unpublishedSequence = null; + yield* Ref.update(historyEpoch, (epoch) => epoch + 1); + yield* SubscriptionRef.set(lastSequence, item.snapshot.snapshotSequence); + yield* setThread(item.snapshot.thread, pageStateFromSnapshot(item.snapshot.page)); + continue; + } - const current = yield* SubscriptionRef.get(state); - if (Option.isNone(current.data)) { - if (item.event.type === "thread.deleted") { + const sequence = unpublishedSequence ?? (yield* SubscriptionRef.get(lastSequence)); + if (item.event.sequence <= sequence) { + continue; + } + unpublishedSequence = item.event.sequence; + + const base = unpublishedThread ?? Option.getOrNull((yield* SubscriptionRef.get(state)).data); + if (base === null) { + if (item.event.type === "thread.deleted") { + yield* publish; + yield* setDeleted(); + } + continue; + } + if (item.event.type === "thread.reverted") { + // A revert rewrites loaded history (whole turns disappear), so an + // older-page fetch in flight may straddle the removed range; the epoch + // bump discards it. The stored page cursor stays valid: cursors are an + // (anchor, turnId) keyset derived from event content, which survives + // the revert projector's row rewrite, so no refresh is needed — the + // revert reducer's turn filtering fully handles loaded history. + yield* Ref.update(historyEpoch, (epoch) => epoch + 1); + } + const result = applyThreadDetailEvent(base, item.event); + if (result.kind === "updated") { + unpublishedThread = result.thread; + } else if (result.kind === "deleted") { + unpublishedThread = null; + yield* publish; yield* setDeleted(); } - return; - } - if (item.event.type === "thread.reverted") { - // A revert rewrites loaded history (whole turns disappear), so an - // older-page fetch in flight may straddle the removed range; the epoch - // bump discards it. The stored page cursor stays valid: cursors are an - // (anchor, turnId) keyset derived from event content, which survives - // the revert projector's row rewrite, so no refresh is needed — the - // revert reducer's turn filtering fully handles loaded history. - yield* Ref.update(historyEpoch, (epoch) => epoch + 1); - } - const result = applyThreadDetailEvent(current.data.value, item.event); - if (result.kind === "updated") { - yield* setThread(result.thread, "keep"); - } else if (result.kind === "deleted") { - yield* setDeleted(); + // The event may have advanced the live state past a parked page's + // watermark; merge it as soon as that happens, before a later event + // in the batch could touch the rows the page carries. + if ((yield* Ref.get(pendingOlderPage)) !== null) { + yield* publish; + yield* tryMergePendingOlderPage(); + } } - // The event may have advanced the live state past a parked page's - // watermark; merge it as soon as that happens. - yield* tryMergePendingOlderPage(); + yield* publish; }); // Merges a parked older page once the live state has caught up to the @@ -399,11 +430,9 @@ export const makeEnvironmentThreadState = Effect.fn("EnvironmentThreadState.make }, ); - const applyItem = Effect.fn("EnvironmentThreadState.applyItem")(function* ( - item: OrchestrationThreadStreamItem, - ) { - yield* applyLock.withPermits(1)(applyItemLocked(item)); - }); + const applyItems = (items: ReadonlyArray) => + applyLock.withPermits(1)(applyItemsLocked(items)); + const applyItem = (item: OrchestrationThreadStreamItem) => applyItems([item]); // Merges an older disjoint page below the currently loaded window. All four // windowed collections prepend; identity dedupe guards the (server-bug or @@ -641,7 +670,12 @@ export const makeEnvironmentThreadState = Effect.fn("EnvironmentThreadState.make retryExpectedFailureAfter: "250 millis", resubscribe: foregroundResubscriptions, }, - ).pipe(Stream.runForEach(applyItem)), + ).pipe( + // Items the transport delivered together (one socket read, or a backlog + // queued while the previous batch applied) arrive as one array and + // publish once. + Stream.runForEachArray(applyItems), + ), ); // Expose loadOlderTurns to UI actions through the request registry.