diff --git a/apps/mobile/src/components/code-reviewer/review-detail-screen.mounted.test.tsx b/apps/mobile/src/components/code-reviewer/review-detail-screen.mounted.test.tsx index 6e67a85650..c54dfe63b5 100644 --- a/apps/mobile/src/components/code-reviewer/review-detail-screen.mounted.test.tsx +++ b/apps/mobile/src/components/code-reviewer/review-detail-screen.mounted.test.tsx @@ -54,6 +54,10 @@ const spectatorQueries = vi.hoisted(() => ({ data: null as unknown, refetch: vi.fn(), }, + sessionMessagesQuery: null as { + enabled?: boolean; + refetchInterval?: unknown; + } | null, })); const statusHelpers = vi.hoisted(() => ({ cancellable: false, @@ -170,11 +174,13 @@ vi.mock('@/lib/trpc', () => ({ }), })); vi.mock('@tanstack/react-query', () => ({ - useQuery: (options: { queryKey?: unknown[] }) => { + useQuery: (options: { queryKey?: unknown[]; enabled?: boolean; refetchInterval?: unknown }) => { const key = options.queryKey?.[0]; - return key === 'codeReviews.getReviewStreamInfo' - ? spectatorQueries.streamInfo - : spectatorQueries.sessionMessages; + if (key === 'codeReviews.getReviewStreamInfo') { + return spectatorQueries.streamInfo; + } + spectatorQueries.sessionMessagesQuery = options; + return spectatorQueries.sessionMessages; }, })); vi.mock('@/components/code-reviewer/review-spectator-stream', () => ({ @@ -307,6 +313,7 @@ beforeEach(() => { spectatorQueries.sessionMessages.isError = false; spectatorQueries.sessionMessages.data = { success: true, entries: [] }; spectatorQueries.sessionMessages.refetch.mockClear(); + spectatorQueries.sessionMessagesQuery = null; spectatorStream.createReviewSpectatorStream.mockReset(); spectatorStream.createReviewSpectatorStream.mockResolvedValue({ connect: vi.fn(), @@ -666,6 +673,57 @@ describe('ReviewDetailScreen spectator transcript', () => { expect(texts).toContain('Waiting for the review transcript.'); }); + it('polls session messages for an in-progress org review and does not open a websocket', () => { + spectatorQueries.streamInfo.data = makeStreamInfo({ + status: 'running', + cloudAgentSessionId: 'agent-1', + organizationId: 'org-1', + }); + spectatorQueries.sessionMessages.data = { + success: true, + entries: [{ timestamp: 't1', message: 'Tool: read', eventType: 'tool' }], + }; + detail.data = { + success: true, + review: makeReview({ status: 'running' }), + tokenUsage: { input: 0, output: 0 }, + }; + + renderScreen(true); + + expect(spectatorStream.createReviewSpectatorStream).not.toHaveBeenCalled(); + expect(spectatorQueries.sessionMessagesQuery?.enabled).toBe(true); + expect(spectatorQueries.sessionMessagesQuery?.refetchInterval).toBe(2000); + const items = sessionListRenders.list.at(-1)?.items as { message?: string }[] | undefined; + expect(items?.map(item => item.message)).toEqual(['Tool: read']); + }); + + it('keeps org poll rows when a later snapshot is empty', () => { + spectatorQueries.streamInfo.data = makeStreamInfo({ + status: 'running', + cloudAgentSessionId: 'agent-1', + organizationId: 'org-1', + }); + spectatorQueries.sessionMessages.data = { + success: true, + entries: [{ timestamp: 't1', message: 'Tool: read', eventType: 'tool' }], + }; + detail.data = { + success: true, + review: makeReview({ status: 'running' }), + tokenUsage: { input: 0, output: 0 }, + }; + + const renderer = mountScreen(true); + spectatorQueries.sessionMessages.data = { success: true, entries: [] }; + act(() => { + renderer.update(createElement(ReviewDetailScreen, { scope: 'personal', reviewId: 'rev-1' })); + }); + + const items = sessionListRenders.list.at(-1)?.items as { message?: string }[] | undefined; + expect(items?.map(item => item.message)).toEqual(['Tool: read']); + }); + it('shows empty copy for a completed review without a session', () => { spectatorQueries.streamInfo.data = makeStreamInfo({ status: 'completed' }); detail.data = { diff --git a/apps/mobile/src/components/code-reviewer/review-spectator-behavior.test.ts b/apps/mobile/src/components/code-reviewer/review-spectator-behavior.test.ts new file mode 100644 index 0000000000..58493a22bf --- /dev/null +++ b/apps/mobile/src/components/code-reviewer/review-spectator-behavior.test.ts @@ -0,0 +1,181 @@ +import { describe, expect, it } from 'vitest'; + +import { + getCodeReviewDisplayBehavior, + resolveReviewSpectatorMode, + retainPolledSpectatorRows, + reviewSpectatorStreamInfoInterval, +} from './review-spectator-behavior'; + +describe('getCodeReviewDisplayBehavior', () => { + it('loads persisted history without polling for a nonterminal V1 review', () => { + expect( + getCodeReviewDisplayBehavior({ + agentVersion: 'v1', + status: 'running', + }) + ).toEqual({ + isHistorical: true, + isTerminal: false, + shouldLoadMessages: true, + shouldPollMessages: false, + shouldPollStatus: false, + }); + }); + + it('keeps a personal V2 review on the live stream path while polling its status', () => { + expect( + getCodeReviewDisplayBehavior({ + agentVersion: 'v2', + status: 'running', + }) + ).toEqual({ + isHistorical: false, + isTerminal: false, + shouldLoadMessages: false, + shouldPollMessages: false, + shouldPollStatus: true, + }); + }); + + it.each(['pending', 'queued', 'running'])( + 'polls organization review transcripts when %s', + status => { + expect( + getCodeReviewDisplayBehavior({ + agentVersion: 'v2', + status, + organizationId: 'org-1', + }) + ).toEqual({ + isHistorical: false, + isTerminal: false, + shouldLoadMessages: true, + shouldPollMessages: true, + shouldPollStatus: true, + }); + } + ); + + it.each(['completed', 'failed', 'cancelled', 'interrupted'])( + 'loads the transcript without polling when %s', + status => { + expect( + getCodeReviewDisplayBehavior({ + agentVersion: 'v2', + status, + organizationId: 'org-1', + }) + ).toEqual({ + isHistorical: false, + isTerminal: true, + shouldLoadMessages: true, + shouldPollMessages: false, + shouldPollStatus: false, + }); + } + ); +}); + +describe('retainPolledSpectatorRows', () => { + it('keeps the last non-empty poll when the latest snapshot is empty', () => { + expect(retainPolledSpectatorRows([], ['kept'], true)).toEqual(['kept']); + }); + + it('uses the latest snapshot when it has rows', () => { + expect(retainPolledSpectatorRows(['next'], ['kept'], true)).toEqual(['next']); + }); + + it('does not retain empty history after polling stops', () => { + expect(retainPolledSpectatorRows([], ['kept'], false)).toEqual([]); + }); +}); + +describe('reviewSpectatorStreamInfoInterval', () => { + it('polls while stream info has not loaded', () => { + expect(reviewSpectatorStreamInfoInterval(undefined)).toBe(2000); + }); + + it('polls an in-flight v2 review even after the session id appears', () => { + expect( + reviewSpectatorStreamInfoInterval({ + success: true, + agentVersion: 'v2', + status: 'running', + cloudAgentSessionId: 'agent-1', + }) + ).toBe(2000); + }); + + it('stops polling a terminal review', () => { + expect( + reviewSpectatorStreamInfoInterval({ + success: true, + agentVersion: 'v2', + status: 'completed', + cloudAgentSessionId: 'agent-1', + }) + ).toBe(false); + }); +}); + +describe('resolveReviewSpectatorMode', () => { + it('polls messages for an in-progress org review and does not open a live stream', () => { + expect( + resolveReviewSpectatorMode( + { + agentVersion: 'v2', + status: 'running', + organizationId: 'org-1', + cloudAgentSessionId: 'agent-1', + }, + 'running', + 0 + ) + ).toEqual({ + isTerminal: false, + shouldPollMessages: true, + shouldLoadHistory: true, + liveCloudId: null, + }); + }); + + it('opens a live stream for an in-progress personal review', () => { + expect( + resolveReviewSpectatorMode( + { + agentVersion: 'v2', + status: 'running', + cloudAgentSessionId: 'agent-1', + }, + 'running', + 0 + ) + ).toEqual({ + isTerminal: false, + shouldPollMessages: false, + shouldLoadHistory: false, + liveCloudId: 'agent-1', + }); + }); + + it('loads history when the parent status is already terminal', () => { + expect( + resolveReviewSpectatorMode( + { + agentVersion: 'v2', + status: 'running', + organizationId: 'org-1', + cloudAgentSessionId: 'agent-1', + }, + 'completed', + 0 + ) + ).toMatchObject({ + isTerminal: true, + shouldPollMessages: false, + shouldLoadHistory: true, + liveCloudId: null, + }); + }); +}); diff --git a/apps/mobile/src/components/code-reviewer/review-spectator-behavior.ts b/apps/mobile/src/components/code-reviewer/review-spectator-behavior.ts new file mode 100644 index 0000000000..a771f68190 --- /dev/null +++ b/apps/mobile/src/components/code-reviewer/review-spectator-behavior.ts @@ -0,0 +1,108 @@ +const TERMINAL_REVIEW_STATUSES = new Set(['completed', 'failed', 'cancelled', 'interrupted']); + +type CodeReviewStreamSnapshot = { + agentVersion: string; + status: string; + organizationId?: string; +}; + +type CodeReviewDisplayBehavior = { + isHistorical: boolean; + isTerminal: boolean; + shouldLoadMessages: boolean; + shouldPollMessages: boolean; + shouldPollStatus: boolean; +}; + +/** + * Same gates as apps/web `getCodeReviewDisplayBehavior`. Org reviews run as + * bot-owned sessions, so a stream ticket is creator-only (web #5781). In-flight + * org transcripts must poll `getSessionMessages` instead of opening a socket. + */ +export function getCodeReviewDisplayBehavior( + snapshot: CodeReviewStreamSnapshot +): CodeReviewDisplayBehavior { + const isHistorical = snapshot.agentVersion !== 'v2'; + const isTerminal = TERMINAL_REVIEW_STATUSES.has(snapshot.status); + const shouldPollStatus = !isHistorical && !isTerminal; + const shouldPollMessages = shouldPollStatus && Boolean(snapshot.organizationId); + + return { + isHistorical, + isTerminal, + shouldLoadMessages: isHistorical || isTerminal || shouldPollMessages, + shouldPollMessages, + shouldPollStatus, + }; +} + +/** Keep the last non-empty poll so an empty ingest snapshot cannot blank the log. */ +export function retainPolledSpectatorRows( + latest: readonly T[], + retained: readonly T[], + shouldPoll: boolean +): readonly T[] { + if (shouldPoll && latest.length === 0 && retained.length > 0) { + return retained; + } + return latest; +} + +type StreamInfo = { + agentVersion: string; + status: string; + organizationId?: string; + cloudAgentSessionId: string | null; +}; + +export function reviewSpectatorStreamInfoInterval( + data: ({ success?: boolean } & Partial) | undefined +): number | false { + if (!data?.success || data.agentVersion === undefined || data.status === undefined) { + return 2000; + } + return getCodeReviewDisplayBehavior({ + agentVersion: data.agentVersion, + status: data.status, + organizationId: data.organizationId, + }).shouldPollStatus + ? 2000 + : false; +} + +type ReviewSpectatorMode = { + isTerminal: boolean; + shouldPollMessages: boolean; + shouldLoadHistory: boolean; + liveCloudId: string | null; +}; + +export function resolveReviewSpectatorMode( + info: StreamInfo | null, + parentStatus: string, + liveRowCount: number +): ReviewSpectatorMode { + const parentIsTerminal = TERMINAL_REVIEW_STATUSES.has(parentStatus); + if (info === null) { + return { + isTerminal: parentIsTerminal, + shouldPollMessages: false, + shouldLoadHistory: false, + liveCloudId: null, + }; + } + const displayBehavior = getCodeReviewDisplayBehavior({ + agentVersion: info.agentVersion, + status: parentIsTerminal ? parentStatus : info.status, + organizationId: info.organizationId, + }); + return { + isTerminal: displayBehavior.isTerminal, + shouldPollMessages: displayBehavior.shouldPollMessages, + shouldLoadHistory: displayBehavior.shouldLoadMessages && liveRowCount === 0, + liveCloudId: + displayBehavior.shouldLoadMessages || info.cloudAgentSessionId === null + ? null + : info.cloudAgentSessionId, + }; +} diff --git a/apps/mobile/src/components/code-reviewer/review-spectator-live.ts b/apps/mobile/src/components/code-reviewer/review-spectator-live.ts new file mode 100644 index 0000000000..3b06a0dbd7 --- /dev/null +++ b/apps/mobile/src/components/code-reviewer/review-spectator-live.ts @@ -0,0 +1,102 @@ +import { type Dispatch, type SetStateAction, useEffect } from 'react'; +import { type TFunction } from 'i18next'; + +import { + appendSpectatorRows, + createSpectatorRowBatcher, + type SpectatorRow, + toSpectatorRow, +} from '@/components/code-reviewer/review-spectator-rows'; +import { + type Connection, + createReviewSpectatorStream, +} from '@/components/code-reviewer/review-spectator-stream'; + +export function useReviewSpectatorLiveStream(input: { + liveCloudId: string | null; + organizationId?: string; + retryNonce: number; + t: TFunction; + setLiveRows: Dispatch>; + setLiveError: Dispatch>; +}): void { + const { liveCloudId, organizationId, retryNonce, t, setLiveRows, setLiveError } = input; + + useEffect(() => { + // Each effect run owns its dispose flag: a shared flag is reset at entry by + // the next run, so a superseded start would see `false` after its own + // cleanup ran and leave a second live socket. A run-local flag keeps the + // stale start from calling `connect()` and makes it destroy its connection. + let disposed = false; + let connection: Connection | null = null; + const clearLiveError = () => { + if (!disposed) { + setLiveError(false); + } + }; + const batcher = createSpectatorRowBatcher(batch => { + if (!disposed) { + setLiveRows(prev => appendSpectatorRows(prev, batch)); + } + }); + + void (async () => { + if (liveCloudId === null) { + return; + } + setLiveError(false); + try { + const created = await createReviewSpectatorStream({ + cloudAgentSessionId: liveCloudId, + organizationId, + onEvent: event => { + if (disposed) { + return; + } + const row = toSpectatorRow(event, t); + if (row === null) { + return; + } + // Synthetic events (connected, snapshots, queued messages) all carry + // eventId 0. A shared key would collapse them into one row. + const keyedRow = + row.key === undefined && event.eventId > 0 + ? { ...row, key: `event-${event.eventId}` } + : row; + batcher.push(keyedRow); + }, + onConnected: clearLiveError, + onReconnected: clearLiveError, + onDisconnected: () => { + if (!disposed) { + setLiveError(true); + } + }, + onError: () => { + if (!disposed) { + setLiveError(true); + } + }, + }); + // oxlint-disable-next-line typescript-eslint/no-unnecessary-condition -- the cleanup closure sets `disposed` after this await resolves + if (disposed) { + created.destroy(); + return; + } + connection = created; + created.connect(); + } catch { + // oxlint-disable-next-line typescript-eslint/no-unnecessary-condition -- the cleanup closure sets `disposed` before this catch can run + if (!disposed) { + setLiveError(true); + } + } + })(); + + return () => { + disposed = true; + batcher.dispose(); + connection?.destroy(); + }; + }, [liveCloudId, organizationId, retryNonce, t, setLiveRows, setLiveError]); +} diff --git a/apps/mobile/src/components/code-reviewer/review-spectator.tsx b/apps/mobile/src/components/code-reviewer/review-spectator.tsx index b836653eb7..cfe7ca2618 100644 --- a/apps/mobile/src/components/code-reviewer/review-spectator.tsx +++ b/apps/mobile/src/components/code-reviewer/review-spectator.tsx @@ -1,5 +1,5 @@ import { type ListRenderItem } from '@shopify/flash-list'; -import { useEffect, useMemo, useState } from 'react'; +import { useMemo, useRef, useState } from 'react'; import { useTranslation } from 'react-i18next'; import { View } from 'react-native'; import { useSafeAreaInsets } from 'react-native-safe-area-context'; @@ -7,27 +7,24 @@ import { useQuery } from '@tanstack/react-query'; import { SessionMessageList } from '@/components/agents/session-message-list'; import { SessionSkeletonMessages } from '@/components/agents/session-detail-skeleton'; +import { + resolveReviewSpectatorMode, + retainPolledSpectatorRows, + reviewSpectatorStreamInfoInterval, +} from '@/components/code-reviewer/review-spectator-behavior'; +import { useReviewSpectatorLiveStream } from '@/components/code-reviewer/review-spectator-live'; import { CompactRetry } from '@/components/code-reviewer/review-spectator-retry'; import { - appendSpectatorRows, - createSpectatorRowBatcher, formatSpectatorTime, type SpectatorRow, spectatorRowsFromEntries, - toSpectatorRow, } from '@/components/code-reviewer/review-spectator-rows'; -import { - type Connection, - createReviewSpectatorStream, -} from '@/components/code-reviewer/review-spectator-stream'; import { useRefetchSessionMessagesOnTerminal } from '@/components/code-reviewer/review-spectator-terminal-refetch'; import { CenteredState } from '@/components/centered-state'; import { QueryError } from '@/components/query-error'; import { Text } from '@/components/ui/text'; import { useTRPC } from '@/lib/trpc'; -const TERMINAL_REVIEW_STATUSES = new Set(['completed', 'failed', 'cancelled', 'interrupted']); - const renderSpectatorRow: ListRenderItem = ({ item }) => ( @@ -74,20 +71,7 @@ export function ReviewSpectator({ const streamInfo = useQuery({ ...trpc.codeReviews.getReviewStreamInfo.queryOptions({ reviewId }), - refetchInterval: query => { - const data = query.state.data; - if (!data?.success) { - return 2000; - } - const isTerminal = TERMINAL_REVIEW_STATUSES.has(data.status); - const isHistorical = data.agentVersion !== 'v2'; - // Poll only while an in-flight v2 review has not yet exposed its - // cloud-agent session; everything else is stable. - if (!isTerminal && !isHistorical && !data.cloudAgentSessionId) { - return 2000; - } - return false; - }, + refetchInterval: query => reviewSpectatorStreamInfoInterval(query.state.data), }); const [liveRows, setLiveRows] = useState([]); @@ -95,109 +79,43 @@ export function ReviewSpectator({ const [retryNonce, setRetryNonce] = useState(0); const info = streamInfo.data?.success ? streamInfo.data : null; - const effectiveStatus = info?.status ?? status; - const isTerminal = - TERMINAL_REVIEW_STATUSES.has(effectiveStatus) || TERMINAL_REVIEW_STATUSES.has(status); - const isHistorical = info !== null && info.agentVersion !== 'v2'; - const shouldLoadHistory = info !== null && (isHistorical || isTerminal) && liveRows.length === 0; - const cloudAgentSessionId = info?.cloudAgentSessionId ?? null; - const isLiveCloud = - info !== null && !isTerminal && info.agentVersion === 'v2' && cloudAgentSessionId !== null; - const liveCloudId = isLiveCloud ? cloudAgentSessionId : null; + const { isTerminal, shouldPollMessages, shouldLoadHistory, liveCloudId } = + resolveReviewSpectatorMode(info, status, liveRows.length); + const isLiveCloud = liveCloudId !== null; const sessionMessages = useQuery({ ...trpc.codeReviews.getSessionMessages.queryOptions({ reviewId }), enabled: Boolean(info) && shouldLoadHistory, + refetchInterval: shouldPollMessages ? 2000 : false, }); useRefetchSessionMessagesOnTerminal(isTerminal, shouldLoadHistory, sessionMessages.refetch); - - useEffect(() => { - // Each effect run owns its dispose flag: a shared flag is reset at entry by - // the next run, so a superseded start would see `false` after its own - // cleanup ran and leave a second live socket. A run-local flag keeps the - // stale start from calling `connect()` and makes it destroy its connection. - let disposed = false; - let connection: Connection | null = null; - const clearLiveError = () => { - if (!disposed) { - setLiveError(false); - } - }; - const batcher = createSpectatorRowBatcher(batch => { - if (!disposed) { - setLiveRows(prev => appendSpectatorRows(prev, batch)); - } - }); - - void (async () => { - if (liveCloudId === null) { - return; - } - setLiveError(false); - try { - const created = await createReviewSpectatorStream({ - cloudAgentSessionId: liveCloudId, - organizationId: info?.organizationId, - onEvent: event => { - if (disposed) { - return; - } - const row = toSpectatorRow(event, t); - if (row === null) { - return; - } - // Synthetic events (connected, snapshots, queued messages) all carry - // eventId 0. A shared key would collapse them into one row. - const keyedRow = - row.key === undefined && event.eventId > 0 - ? { ...row, key: `event-${event.eventId}` } - : row; - batcher.push(keyedRow); - }, - onConnected: clearLiveError, - onReconnected: clearLiveError, - onDisconnected: () => { - if (!disposed) { - setLiveError(true); - } - }, - onError: () => { - if (!disposed) { - setLiveError(true); - } - }, - }); - // oxlint-disable-next-line typescript-eslint/no-unnecessary-condition -- the cleanup closure sets `disposed` after this await resolves - if (disposed) { - created.destroy(); - return; - } - connection = created; - created.connect(); - } catch { - // oxlint-disable-next-line typescript-eslint/no-unnecessary-condition -- the cleanup closure sets `disposed` before this catch can run - if (!disposed) { - setLiveError(true); - } - } - })(); - - return () => { - disposed = true; - batcher.dispose(); - connection?.destroy(); - }; - }, [liveCloudId, info?.organizationId, retryNonce, t]); + useReviewSpectatorLiveStream({ + liveCloudId, + organizationId: info?.organizationId, + retryNonce, + t, + setLiveRows, + setLiveError, + }); const historicalRows = useMemo( () => sessionMessages.data?.success ? spectatorRowsFromEntries(sessionMessages.data.entries) : [], [sessionMessages.data] ); + const retainedHistoryRef = useRef([]); + const displayHistory = retainPolledSpectatorRows( + historicalRows, + retainedHistoryRef.current, + shouldPollMessages + ); + if (historicalRows.length > 0) { + retainedHistoryRef.current = historicalRows; + } const transcriptRows = - shouldLoadHistory || (info === null && liveRows.length === 0) ? historicalRows : liveRows; + shouldLoadHistory || (info === null && liveRows.length === 0) ? displayHistory : liveRows; function renderRowsWithRetry(onRetry: () => void) { return (