feat(desktop): paginate stream particles via windowed Firestore subscription #296

Open
talksik wants to merge 2 commits from claude/paginated-firestorm-particles-omb5bs into main
7 changed files with 464 additions and 94 deletions
+6 -1
View File
@@ -6,8 +6,12 @@ import { cn } from '@/lib/utils';
function ScrollArea({ function ScrollArea({
className, className,
children, children,
viewportRef,
...props ...props
}: React.ComponentProps<typeof ScrollAreaPrimitive.Root>) { }: React.ComponentProps<typeof ScrollAreaPrimitive.Root> & {
/** Ref to the scrollable viewport, e.g. for scroll anchoring or observers. */
viewportRef?: React.Ref<HTMLDivElement>;
}) {
return ( return (
<ScrollAreaPrimitive.Root <ScrollAreaPrimitive.Root
data-slot="scroll-area" data-slot="scroll-area"
@@ -15,6 +19,7 @@ function ScrollArea({
{...props} {...props}
> >
<ScrollAreaPrimitive.Viewport <ScrollAreaPrimitive.Viewport
ref={viewportRef}
data-slot="scroll-area-viewport" data-slot="scroll-area-viewport"
className="focus-visible:ring-ring/50 size-full rounded-[inherit] transition-[color,box-shadow] outline-none focus-visible:ring-[3px] focus-visible:outline-1 [&>div]:!w-full" className="focus-visible:ring-ring/50 size-full rounded-[inherit] transition-[color,box-shadow] outline-none focus-visible:ring-[3px] focus-visible:outline-1 [&>div]:!w-full"
> >
@@ -7,58 +7,87 @@ import {
import type { HumanPresence } from '@/hooks/use-presence-positions'; import type { HumanPresence } from '@/hooks/use-presence-positions';
const MAX_VISIBLE_AVATARS = 3; const MAX_VISIBLE_AVATARS = 3;
const PAGE_SIZE = 10; const VISIBLE_SEGMENTS = 10;
interface PlaybackPageIndicatorProps { interface PlaybackPageIndicatorProps {
total: number; /** Number of particles currently loaded in the window. */
loadedCount: number;
current: number; current: number;
progress: number; progress: number;
onGoTo: (index: number) => void; onGoTo: (index: number) => void;
presenceBySegment?: Map<number, HumanPresence[]>; presenceBySegment?: Map<number, HumanPresence[]>;
/** Set of humanIds currently online in the stream channel. */ /** Set of humanIds currently online in the stream channel. */
onlineHumanIds?: Set<string>; onlineHumanIds?: Set<string>;
/** More (older) particles exist before the loaded window. */
hasMoreOlder?: boolean;
/** Pull in older history when scrubbing past the oldest loaded segment. */
onLoadOlder?: () => void;
/** Render only avatars or only tracks. Omit to render both. */ /** Render only avatars or only tracks. Omit to render both. */
layer?: 'avatars' | 'tracks'; layer?: 'avatars' | 'tracks';
} }
/**
* A streaming, count-free progress scrubber. The stream is paginated, so the
* true particle count is unknown — instead this shows a sliding window of
* segments around the current position. Segments map to the loaded particles
* (newest on the right); the edge stubs scrub within the window and pull in
* older history when you reach the oldest loaded segment.
*/
export function PlaybackPageIndicator({ export function PlaybackPageIndicator({
total, loadedCount,
current, current,
progress, progress,
onGoTo, onGoTo,
presenceBySegment, presenceBySegment,
onlineHumanIds, onlineHumanIds,
hasMoreOlder,
onLoadOlder,
layer, layer,
}: PlaybackPageIndicatorProps) { }: PlaybackPageIndicatorProps) {
if (total === 0) return null; if (loadedCount === 0) return null;
const showAvatars = layer !== 'tracks'; const showAvatars = layer !== 'tracks';
const showTracks = layer !== 'avatars'; const showTracks = layer !== 'avatars';
const paginated = total > PAGE_SIZE;
const safeCurrent = current < 0 ? 0 : current; const safeCurrent = current < 0 ? 0 : current;
const pageStart = paginated // Slide the visible window so the current segment stays in view with a bit of
? Math.floor(safeCurrent / PAGE_SIZE) * PAGE_SIZE // context on either side, clamped to the loaded range.
: 0; const sliceStart = Math.min(
const visibleCount = paginated Math.max(0, safeCurrent - Math.floor(VISIBLE_SEGMENTS / 2)),
? Math.min(PAGE_SIZE, total - pageStart) Math.max(0, loadedCount - VISIBLE_SEGMENTS),
: total; );
const hasPrevPage = paginated && pageStart > 0; const sliceEnd = Math.min(loadedCount, sliceStart + VISIBLE_SEGMENTS);
const hasNextPage = paginated && pageStart + PAGE_SIZE < total; const visibleCount = sliceEnd - sliceStart;
const paginated = loadedCount > VISIBLE_SEGMENTS || !!hasMoreOlder;
// Older = lower indices (left); newer = higher indices (right).
const hasOlder = sliceStart > 0 || !!hasMoreOlder;
const hasNewer = sliceEnd < loadedCount;
// Stubs jump a page at a time; reaching the oldest loaded pulls in history.
const goOlder = () => {
const target = safeCurrent - VISIBLE_SEGMENTS;
if (target >= 0) onGoTo(target);
else if (sliceStart > 0) onGoTo(0);
else if (hasMoreOlder) onLoadOlder?.();
};
const goNewer = () => {
onGoTo(Math.min(loadedCount - 1, safeCurrent + VISIBLE_SEGMENTS));
};
return ( return (
<div className="flex w-full flex-col items-stretch leading-none"> <div className="flex w-full flex-col items-stretch leading-none">
<div className="flex w-full items-end gap-px"> <div className="flex w-full items-end gap-px">
{paginated && ( {paginated && (
<GhostStub <GhostStub
visible={hasPrevPage} visible={hasOlder}
interactive={showTracks} interactive={showTracks}
onClick={() => onGoTo(pageStart - 1)} onClick={goOlder}
/> />
)} )}
<div className="flex flex-1 items-end gap-px"> <div className="flex flex-1 items-end gap-px">
{Array.from({ length: visibleCount }, (_, j) => { {Array.from({ length: visibleCount }, (_, j) => {
const i = pageStart + j; const i = sliceStart + j;
const presence = presenceBySegment?.get(i); const presence = presenceBySegment?.get(i);
return ( return (
<div key={i} className="flex flex-1 flex-col items-stretch"> <div key={i} className="flex flex-1 flex-col items-stretch">
@@ -100,17 +129,12 @@ export function PlaybackPageIndicator({
</div> </div>
{paginated && ( {paginated && (
<GhostStub <GhostStub
visible={hasNextPage} visible={hasNewer}
interactive={showTracks} interactive={showTracks}
onClick={() => onGoTo(pageStart + PAGE_SIZE)} onClick={goNewer}
/> />
)} )}
</div> </div>
{paginated && showTracks && current >= 0 && (
<div className="pointer-events-none pt-1 text-center text-[10px] font-medium tabular-nums tracking-wide text-white/40">
{current + 1} / {total}
</div>
)}
</div> </div>
); );
} }
@@ -7,24 +7,28 @@ import type { HumanPresence } from '@/hooks/use-presence-positions';
export function BottomBar({ export function BottomBar({
visible, visible,
total, loadedCount,
current, current,
progress, progress,
onGoTo, onGoTo,
presenceBySegment, presenceBySegment,
onlineHumanIds, onlineHumanIds,
hasMoreOlder,
onLoadOlder,
exitRemainingMs, exitRemainingMs,
onOpenKeybindings, onOpenKeybindings,
onOpenHuddle, onOpenHuddle,
onExit, onExit,
}: { }: {
visible: boolean; visible: boolean;
total: number; loadedCount: number;
current: number; current: number;
progress: number; progress: number;
onGoTo: (index: number) => void; onGoTo: (index: number) => void;
presenceBySegment: Map<number, HumanPresence[]>; presenceBySegment: Map<number, HumanPresence[]>;
onlineHumanIds: Set<string>; onlineHumanIds: Set<string>;
hasMoreOlder: boolean;
onLoadOlder: () => void;
exitRemainingMs: number | null; exitRemainingMs: number | null;
onOpenKeybindings: () => void; onOpenKeybindings: () => void;
onOpenHuddle: () => void; onOpenHuddle: () => void;
@@ -41,21 +45,25 @@ export function BottomBar({
> >
{/* Presence avatars — above the blurred background */} {/* Presence avatars — above the blurred background */}
<PlaybackPageIndicator <PlaybackPageIndicator
total={total} loadedCount={loadedCount}
current={current} current={current}
progress={progress} progress={progress}
onGoTo={onGoTo} onGoTo={onGoTo}
presenceBySegment={presenceBySegment} presenceBySegment={presenceBySegment}
onlineHumanIds={onlineHumanIds} onlineHumanIds={onlineHumanIds}
hasMoreOlder={hasMoreOlder}
onLoadOlder={onLoadOlder}
layer="avatars" layer="avatars"
/> />
{/* Blurred background container — tracks + controls */} {/* Blurred background container — tracks + controls */}
<div className="pb-3"> <div className="pb-3">
<PlaybackPageIndicator <PlaybackPageIndicator
total={total} loadedCount={loadedCount}
current={current} current={current}
progress={progress} progress={progress}
onGoTo={onGoTo} onGoTo={onGoTo}
hasMoreOlder={hasMoreOlder}
onLoadOlder={onLoadOlder}
layer="tracks" layer="tracks"
/> />
<div className="flex items-center justify-center px-3 pt-2 gap-2"> <div className="flex items-center justify-center px-3 pt-2 gap-2">
@@ -1,5 +1,13 @@
import { useEffect, useRef } from 'react'; import { useCallback, useEffect, useLayoutEffect, useRef } from 'react';
import { CircleCheck, FileText, Image, List, Mic, Video } from 'lucide-react'; import {
CircleCheck,
FileText,
Image,
List,
Loader2,
Mic,
Video,
} from 'lucide-react';
import { isParticleDeleted, type Human, type Particle } from '@/api/types'; import { isParticleDeleted, type Human, type Particle } from '@/api/types';
import { cn } from '@/lib/utils'; import { cn } from '@/lib/utils';
import { useNetwork } from '@/hooks/use-networks'; import { useNetwork } from '@/hooks/use-networks';
@@ -15,12 +23,18 @@ interface StreamListSidebarProps {
currentIndex: number; currentIndex: number;
onSelect: (index: number) => void; onSelect: (index: number) => void;
onToggle: () => void; onToggle: () => void;
/** More (older) particles exist before the loaded window. */
hasMoreOlder: boolean;
/** Load the next page of older particles. */
onLoadOlder: () => void;
isLoadingOlder: boolean;
} }
/** /**
* Browse-mode panel beside the stream: a chat-like timeline of every * Browse-mode panel beside the stream: a chat-like timeline of the loaded
* particle. Selecting a message plays it in the immersive stream view; * particles. Selecting a message plays it in the immersive stream view;
* nothing auto-advances. * nothing auto-advances. Older history is fetched automatically as the user
* scrolls toward the top.
*/ */
export function StreamListSidebar({ export function StreamListSidebar({
items, items,
@@ -28,22 +42,67 @@ export function StreamListSidebar({
currentIndex, currentIndex,
onSelect, onSelect,
onToggle, onToggle,
hasMoreOlder,
onLoadOlder,
isLoadingOlder,
}: StreamListSidebarProps) { }: StreamListSidebarProps) {
const network = useNetwork(networkId); const network = useNetwork(networkId);
const rowRefs = useRef<Array<HTMLDivElement | null>>([]); const rowRefs = useRef<Array<HTMLDivElement | null>>([]);
const viewportRef = useRef<HTMLDivElement>(null);
const sentinelRef = useRef<HTMLDivElement>(null);
// Scroll into view only when the selection genuinely changes — not when the
// current index shifts because older particles were prepended.
const selectedId = currentIndex >= 0 ? items[currentIndex]?.id : undefined;
const prevSelectedIdRef = useRef<string | undefined>(undefined);
useEffect(() => { useEffect(() => {
if (currentIndex >= 0) { if (selectedId && selectedId !== prevSelectedIdRef.current) {
rowRefs.current[currentIndex]?.scrollIntoView({ block: 'nearest' }); rowRefs.current[currentIndex]?.scrollIntoView({ block: 'nearest' });
} }
}, [currentIndex]); prevSelectedIdRef.current = selectedId;
}, [selectedId, currentIndex]);
// Anchor the viewport when older particles are prepended so the content the
// user is looking at stays put instead of jumping.
const pendingAnchorRef = useRef<{ height: number; top: number } | null>(null);
const requestOlder = useCallback(() => {
const vp = viewportRef.current;
if (!vp) return;
pendingAnchorRef.current = { height: vp.scrollHeight, top: vp.scrollTop };
onLoadOlder();
}, [onLoadOlder]);
useLayoutEffect(() => {
const vp = viewportRef.current;
const anchor = pendingAnchorRef.current;
if (!vp || !anchor) return;
const delta = vp.scrollHeight - anchor.height;
if (delta > 0) vp.scrollTop = anchor.top + delta;
pendingAnchorRef.current = null;
}, [items]);
// Auto-fetch older history when the top sentinel scrolls into view.
useEffect(() => {
const vp = viewportRef.current;
const sentinel = sentinelRef.current;
if (!vp || !sentinel || !hasMoreOlder) return;
const observer = new IntersectionObserver(
(entries) => {
if (entries[0]?.isIntersecting && !isLoadingOlder) requestOlder();
},
{ root: vp, rootMargin: '120px 0px 0px 0px' },
);
observer.observe(sentinel);
return () => observer.disconnect();
}, [hasMoreOlder, isLoadingOlder, requestOlder]);
return ( return (
<aside className="dark flex w-60 shrink-0 flex-col overflow-hidden border-l border-white/10 bg-zinc-950"> <aside className="dark flex w-60 shrink-0 flex-col overflow-hidden border-l border-white/10 bg-zinc-950">
<div className="flex shrink-0 items-center gap-2 border-b border-white/10 px-4 py-3"> <div className="flex shrink-0 items-center gap-2 border-b border-white/10 px-4 py-3">
<List className="size-3.5 text-white/40" /> <List className="size-3.5 text-white/40" />
<span className="truncate text-sm font-medium text-white/90 mr-auto"> <span className="truncate text-sm font-medium text-white/90 mr-auto">
{items.length} messages {items.length}
{hasMoreOlder ? '+' : ''} messages
</span> </span>
<KeyHint <KeyHint
keys="L" keys="L"
@@ -54,8 +113,14 @@ export function StreamListSidebar({
to close to close
</KeyHint> </KeyHint>
</div> </div>
<ScrollArea className="min-h-0 flex-1"> <ScrollArea viewportRef={viewportRef} className="min-h-0 flex-1">
<div className="flex flex-col gap-0.5 px-2 py-2"> <div className="flex flex-col gap-0.5 px-2 py-2">
<div ref={sentinelRef} aria-hidden />
{hasMoreOlder && (
<div className="flex items-center justify-center py-2 text-white/40">
<Loader2 className="size-3.5 animate-spin" />
</div>
)}
{items.map((item, index) => ( {items.map((item, index) => (
<div <div
key={item.id} key={item.id}
@@ -219,6 +219,9 @@ function StreamViewInner({ path, streamParticle }: StreamViewProps) {
currentParticle, currentParticle,
currentIndex, currentIndex,
status, status,
hasMoreOlder,
loadOlder,
isLoadingOlder,
next, next,
prev, prev,
goTo, goTo,
@@ -545,12 +548,14 @@ function StreamViewInner({ path, streamParticle }: StreamViewProps) {
{/* BottomBar — pinned visible while browsing, mouse-activity in player */} {/* BottomBar — pinned visible while browsing, mouse-activity in player */}
<BottomBar <BottomBar
visible={mode === 'list' || controlsVisible} visible={mode === 'list' || controlsVisible}
total={children.length} loadedCount={children.length}
current={currentIndex} current={currentIndex}
progress={progress} progress={progress}
onGoTo={goTo} onGoTo={goTo}
presenceBySegment={presenceBySegment} presenceBySegment={presenceBySegment}
onlineHumanIds={onlineHumanIds} onlineHumanIds={onlineHumanIds}
hasMoreOlder={hasMoreOlder}
onLoadOlder={loadOlder}
exitRemainingMs={exitRemainingMs} exitRemainingMs={exitRemainingMs}
onOpenKeybindings={() => setShowKeybindings(true)} onOpenKeybindings={() => setShowKeybindings(true)}
onOpenHuddle={handleOpenHuddle} onOpenHuddle={handleOpenHuddle}
@@ -573,6 +578,9 @@ function StreamViewInner({ path, streamParticle }: StreamViewProps) {
currentIndex={currentIndex} currentIndex={currentIndex}
onSelect={goTo} onSelect={goTo}
onToggle={toggleViewMode} onToggle={toggleViewMode}
hasMoreOlder={hasMoreOlder}
onLoadOlder={loadOlder}
isLoadingOlder={isLoadingOlder}
/> />
)} )}
</div> </div>
+102 -56
View File
@@ -8,7 +8,7 @@ import {
} from 'react'; } from 'react';
import { useAuthStore } from '@/stores/auth-store'; import { useAuthStore } from '@/stores/auth-store';
import type { Particle } from '@/api/types'; import type { Particle } from '@/api/types';
import { useLiveParticleChildren } from '@/hooks/use-particle'; import { useWindowedStreamParticles } from '@/hooks/use-windowed-stream-particles';
import { toFirestoreDocPath, type ParticlePath } from '@/lib/particle-path'; import { toFirestoreDocPath, type ParticlePath } from '@/lib/particle-path';
import { updateStreamPlaybackMarker } from '@/lib/firestore-particles'; import { updateStreamPlaybackMarker } from '@/lib/firestore-particles';
@@ -92,6 +92,11 @@ interface UseStreamPlaybackResult {
currentIndex: number; currentIndex: number;
status: PlaybackStatus; status: PlaybackStatus;
initialized: boolean; initialized: boolean;
/** Whether older particles exist before the loaded window (list scroll-up). */
hasMoreOlder: boolean;
/** Extend the loaded window backward. */
loadOlder: () => void;
isLoadingOlder: boolean;
next: () => void; next: () => void;
prev: () => void; prev: () => void;
goTo: (index: number) => void; goTo: (index: number) => void;
@@ -113,50 +118,26 @@ export function useStreamPlayback(
): UseStreamPlaybackResult { ): UseStreamPlaybackResult {
const userId = useAuthStore((s) => s.user?.id); const userId = useAuthStore((s) => s.user?.id);
const [state, dispatch] = useReducer(playbackReducer, initialState); const [state, dispatch] = useReducer(playbackReducer, initialState);
// Read via ref so the onAdded subscription callback stays stable.
const marker = userId
? (streamParticle.playback_markers?.[userId] ?? null)
: null;
// Windowed source: only the tail (plus enough history to cover the marker)
// is loaded, instead of every particle in the stream.
const { children, hasMoreOlder, loadOlder, isLoadingOlder } =
useWindowedStreamParticles(path, { marker });
// Read via ref so the new-particle effect stays cheap to reason about.
const autoAdvanceOnNewRef = useRef(autoAdvanceOnNew); const autoAdvanceOnNewRef = useRef(autoAdvanceOnNew);
useEffect(() => { useEffect(() => {
autoAdvanceOnNewRef.current = autoAdvanceOnNew; autoAdvanceOnNewRef.current = autoAdvanceOnNew;
}, [autoAdvanceOnNew]); }, [autoAdvanceOnNew]);
// Track the stream ID we've initialized for, to reset when navigating between streams
// Track the stream ID we've initialized for, to reset when navigating streams.
const initializedForRef = useRef<string | null>(null); const initializedForRef = useRef<string | null>(null);
// Latest currentIndex for onParticleRemoved, which is passed into
// useLiveParticleChildren. Reading it through a ref keeps the callback stable
// (no re-subscription) and breaks the declaration cycle
// children -> currentIndex -> callback -> children. useEffectEvent can't be
// used here — Effect Events may not be passed to another hook.
const currentIndexRef = useRef(0);
// --- Firestore change callbacks --- // Derive current index and particle from ID.
const onParticleAdded = useCallback((particle: Particle) => {
if (!autoAdvanceOnNewRef.current) return;
dispatch({ type: 'PARTICLE_ADDED', particleId: particle.id });
}, []);
const onParticleRemoved = useCallback(
(removed: Particle, updatedChildren: Particle[]) => {
const fallbackIndex = Math.min(
currentIndexRef.current,
updatedChildren.length - 1,
);
const fallback = updatedChildren[Math.max(0, fallbackIndex)];
dispatch({
type: 'PARTICLE_REMOVED',
removedParticleId: removed.id,
fallbackParticleId: fallback?.id ?? null,
});
},
[],
);
const { children } = useLiveParticleChildren(path, {
orderByField: 'created_at',
orderDirection: 'asc',
onAdded: onParticleAdded,
onRemoved: onParticleRemoved,
});
// Derive current index and particle from ID
const currentIndex = useMemo(() => { const currentIndex = useMemo(() => {
if (!state.currentParticleId) return -1; if (!state.currentParticleId) return -1;
return children.findIndex((c) => c.id === state.currentParticleId); return children.findIndex((c) => c.id === state.currentParticleId);
@@ -164,12 +145,64 @@ export function useStreamPlayback(
const currentParticle = currentIndex !== -1 ? children[currentIndex] : null; const currentParticle = currentIndex !== -1 ? children[currentIndex] : null;
// Keep the ref read by onParticleRemoved in sync with the derived index. // Remember the last index the current particle was actually found at, so a
// removal can fall back to a sensible neighbour even though `currentIndex`
// has already gone to -1 by the time we notice.
const lastValidIndexRef = useRef(0);
useEffect(() => { useEffect(() => {
currentIndexRef.current = currentIndex; if (currentIndex >= 0) lastValidIndexRef.current = currentIndex;
}, [currentIndex]); }, [currentIndex]);
// Fallback init — always sees latest children/state via useEffectEvent // --- New tail particle → resume from end ---
// Derive arrivals from the children tail rather than Firestore change events,
// which can't distinguish a genuine new particle from pagination backfill.
const prevNewestIdRef = useRef<string | null>(null);
useEffect(() => {
if (children.length === 0) {
prevNewestIdRef.current = null;
return;
}
const newestId = children[children.length - 1].id;
const prevNewestId = prevNewestIdRef.current;
prevNewestIdRef.current = newestId;
if (prevNewestId === null || newestId === prevNewestId) return;
if (!autoAdvanceOnNewRef.current) return;
// Resume at the first particle added after where playback ended.
const prevIndex = children.findIndex((c) => c.id === prevNewestId);
const firstNew =
prevIndex >= 0
? (children[prevIndex + 1] ?? children[children.length - 1])
: children[children.length - 1];
dispatch({ type: 'PARTICLE_ADDED', particleId: firstNew.id });
}, [children]);
// --- Current particle removed (deletion) → fall back to a neighbour ---
const prevIdsRef = useRef<Set<string>>(new Set());
useEffect(() => {
const id = state.currentParticleId;
const prevIds = prevIdsRef.current;
const currIds = new Set(children.map((c) => c.id));
prevIdsRef.current = currIds;
if (!id || children.length === 0) return;
if (currIds.has(id)) return;
// Only treat as a removal if it was present before — a not-yet-arrived id
// (e.g. optimistic goToParticle) should wait, not fall back.
if (!prevIds.has(id)) return;
const fallbackIndex = Math.min(
lastValidIndexRef.current,
children.length - 1,
);
const fallback = children[Math.max(0, fallbackIndex)];
dispatch({
type: 'PARTICLE_REMOVED',
removedParticleId: id,
fallbackParticleId: fallback?.id ?? null,
});
}, [children, state.currentParticleId]);
// Fallback init — always sees latest children/state via useEffectEvent.
const initFallback = useEffectEvent(() => { const initFallback = useEffectEvent(() => {
if (state.initialized || children.length === 0) return; if (state.initialized || children.length === 0) return;
initializedForRef.current = streamParticle.id; initializedForRef.current = streamParticle.id;
@@ -178,7 +211,7 @@ export function useStreamPlayback(
// --- Init logic: runs on every children change until initialized --- // --- Init logic: runs on every children change until initialized ---
useEffect(() => { useEffect(() => {
// Reset if we navigated to a different stream // Reset if we navigated to a different stream.
if ( if (
initializedForRef.current !== null && initializedForRef.current !== null &&
initializedForRef.current !== streamParticle.id initializedForRef.current !== streamParticle.id
@@ -186,7 +219,7 @@ export function useStreamPlayback(
initializedForRef.current = null; initializedForRef.current = null;
} }
// Already initialized for this stream // Already initialized for this stream.
if (state.initialized && initializedForRef.current === streamParticle.id) if (state.initialized && initializedForRef.current === streamParticle.id)
return; return;
@@ -195,35 +228,39 @@ export function useStreamPlayback(
const playbackPosition = streamParticle.playback_markers?.[userId ?? '']; const playbackPosition = streamParticle.playback_markers?.[userId ?? ''];
if (!playbackPosition) { if (!playbackPosition) {
// No marker — start from the beginning // No marker — start from the start of the loaded window.
initializedForRef.current = streamParticle.id; initializedForRef.current = streamParticle.id;
dispatch({ type: 'INIT', particleId: children[0].id }); dispatch({ type: 'INIT', particleId: children[0].id });
return; return;
} }
// Try to find the marker's target particle // Resume at the first particle after the marker.
const found = children.find( const found = children.find(
(c) => c.created_at.getTime() > playbackPosition.getTime(), (c) => c.created_at.getTime() > playbackPosition.getTime(),
); );
if (found) { if (found) {
initializedForRef.current = streamParticle.id; initializedForRef.current = streamParticle.id;
dispatch({ type: 'INIT', particleId: found.id }); dispatch({ type: 'INIT', particleId: found.id });
return; return;
} else {
initializedForRef.current = streamParticle.id;
dispatch({ type: 'INIT', particleId: children[children.length - 1].id });
} }
// Marker target not found yet — fall back after timeout // Marker is older than everything loaded so far. If the window is still
const timeout = setTimeout(initFallback, INIT_FALLBACK_TIMEOUT_MS); // growing backward to reach it, wait for more particles to arrive.
return () => clearTimeout(timeout); if (hasMoreOlder) {
const timeout = setTimeout(initFallback, INIT_FALLBACK_TIMEOUT_MS);
return () => clearTimeout(timeout);
}
// Reached the start with no particle after the marker → caught up.
initializedForRef.current = streamParticle.id;
dispatch({ type: 'INIT', particleId: children[children.length - 1].id });
}, [ }, [
children, children,
streamParticle.id, streamParticle.id,
streamParticle.playback_markers, streamParticle.playback_markers,
userId, userId,
state.initialized, state.initialized,
hasMoreOlder,
]); ]);
// --- Persist playback marker (only advance forward, never backwards) --- // --- Persist playback marker (only advance forward, never backwards) ---
@@ -237,7 +274,7 @@ export function useStreamPlayback(
lastPersistedMarkerRef.current ?? lastPersistedMarkerRef.current ??
streamParticle.playback_markers?.[userId]; streamParticle.playback_markers?.[userId];
// Only update if advancing beyond the current marker // Only update if advancing beyond the current marker.
if (existingMarker && currentTime.getTime() <= existingMarker.getTime()) if (existingMarker && currentTime.getTime() <= existingMarker.getTime())
return; return;
@@ -266,12 +303,18 @@ export function useStreamPlayback(
}, [children, currentIndex]); }, [children, currentIndex]);
const prev = useCallback(() => { const prev = useCallback(() => {
if (currentIndex <= 0) return; if (currentIndex < 0) return;
if (currentIndex === 0) {
// At the start of the loaded window — pull in older history so the user
// can keep going back.
if (hasMoreOlder) loadOlder();
return;
}
dispatch({ dispatch({
type: 'SET_PARTICLE', type: 'SET_PARTICLE',
particleId: children[currentIndex - 1].id, particleId: children[currentIndex - 1].id,
}); });
}, [children, currentIndex]); }, [children, currentIndex, hasMoreOlder, loadOlder]);
const goTo = useCallback( const goTo = useCallback(
(index: number) => { (index: number) => {
@@ -295,6 +338,9 @@ export function useStreamPlayback(
currentIndex, currentIndex,
status: state.status, status: state.status,
initialized: state.initialized, initialized: state.initialized,
hasMoreOlder,
loadOlder,
isLoadingOlder,
next, next,
prev, prev,
goTo, goTo,
@@ -0,0 +1,214 @@
import { useCallback, useEffect, useRef, useState } from 'react';
import { subscribeToParticleChildren } from '@/lib/firestore-particles';
import type { Particle } from '@/api/types';
import {
toFirestoreChildrenPath,
type ParticlePath,
} from '@/lib/particle-path';
const DEFAULT_PAGE_SIZE = 30;
interface UseWindowedStreamParticlesParams {
/**
* Resume anchor (the viewer's playback marker). The window grows backward
* until it covers this timestamp so the resume particle is always loaded.
* Captured once per stream — advancing the marker during playback does not
* re-window.
*/
marker?: Date | null;
/** How many particles to add per backward growth step. */
pageSize?: number;
}
export interface UseWindowedStreamParticlesResult {
/**
* Loaded window, ascending (oldest → newest). The newest particle in the
* stream is always present — the window only ever grows backward.
*/
children: Particle[];
isLoading: boolean;
error: Error | null;
/** Whether older particles likely exist before the loaded window. */
hasMoreOlder: boolean;
/** Extend the window backward (older history). No-op when nothing remains. */
loadOlder: () => void;
isLoadingOlder: boolean;
}
/** Per-stream mutable tracking that must survive limit-driven re-subscriptions. */
interface WindowTracking {
path: string | null;
/** Newest created_at (ms) seen — tells new tail particles from backfill. */
newestMs: number | null;
/** Oldest created_at (ms) currently loaded. */
oldestMs: number | null;
/** Backward growth target (the marker, ms), frozen on first capture. */
coverageMs: number | null;
}
/**
* Live, windowed view of a stream's particles.
*
* Instead of subscribing to every child (the old behaviour), this keeps a
* `orderBy(created_at desc) limit(N)` window anchored at the newest particle
* and grows it backward on demand. Because the window is anchored at the tail
* it always contains the most recent particles, so new arrivals stream in and
* forward playback never needs a fetch. The window grows backward to:
* 1. cover the resume marker, so playback can start where the user left off;
* 2. service `loadOlder()` when the list view scrolls up.
*
* Output is reversed to ascending order to match the rest of the playback code.
*/
export function useWindowedStreamParticles(
path: ParticlePath | undefined,
{
marker = null,
pageSize = DEFAULT_PAGE_SIZE,
}: UseWindowedStreamParticlesParams = {},
): UseWindowedStreamParticlesResult {
const [children, setChildren] = useState<Particle[]>([]);
const [isLoading, setIsLoading] = useState(true);
const [error, setError] = useState<Error | null>(null);
const [hasMoreOlder, setHasMoreOlder] = useState(false);
const [isLoadingOlder, setIsLoadingOlder] = useState(false);
const [limit, setLimit] = useState(pageSize);
const collectionPath = path ? toFirestoreChildrenPath(path) : null;
const trackingRef = useRef<WindowTracking>({
path: null,
newestMs: null,
oldestMs: null,
coverageMs: null,
});
// Latest marker, read lazily inside the snapshot callback so a late-resolving
// marker (e.g. auth after first paint) still seeds backward coverage.
const markerRef = useRef(marker);
useEffect(() => {
markerRef.current = marker;
}, [marker]);
// Reset window state when the stream changes (render-phase adjustment — the
// blessed alternative to a reset effect, avoids cascading effect renders).
const [trackedPath, setTrackedPath] = useState(collectionPath);
if (trackedPath !== collectionPath) {
setTrackedPath(collectionPath);
setLimit(pageSize);
setChildren([]);
setIsLoading(true);
setError(null);
setHasMoreOlder(false);
setIsLoadingOlder(false);
}
useEffect(() => {
if (!collectionPath) return;
// Reset per-stream tracking on a genuine stream change, but keep it across
// limit-driven re-subscriptions (newest/coverage must persist).
const tracking = trackingRef.current;
if (tracking.path !== collectionPath) {
tracking.path = collectionPath;
tracking.newestMs = null;
tracking.oldestMs = null;
tracking.coverageMs = null;
}
const unsubscribe = subscribeToParticleChildren(collectionPath, {
orderByField: 'created_at',
orderDirection: 'desc',
limit,
onData: (descData) => {
const t = trackingRef.current;
// Lazily freeze the backward-coverage target from the marker.
if (t.coverageMs === null && markerRef.current) {
t.coverageMs = markerRef.current.getTime();
}
// Firestore caps results at `limit`; a full window means more older
// particles may exist beyond it.
const saturated = descData.length === limit;
const newest = descData[0];
const newestMs = newest ? newest.created_at.getTime() : null;
// Anti-eviction: if the window is full and genuinely newer particles
// arrived at the tail, grow the limit so the oldest loaded particles
// aren't pushed out. Skip rendering the evicted snapshot — the regrown
// query delivers the complete window a beat later.
const prevNewest = t.newestMs;
if (
saturated &&
newestMs !== null &&
prevNewest !== null &&
newestMs > prevNewest
) {
const newerCount = descData.filter(
(d) => d.created_at.getTime() > prevNewest,
).length;
if (newerCount > 0) {
t.newestMs = newestMs;
setLimit((l) => l + newerCount);
return;
}
}
if (newestMs !== null) t.newestMs = newestMs;
const ascData = descData.slice().reverse();
t.oldestMs =
ascData.length > 0 ? ascData[0].created_at.getTime() : null;
// Marker coverage: keep growing backward until the resume marker falls
// within the window (or we reach the start of the stream).
if (
saturated &&
t.coverageMs !== null &&
t.oldestMs !== null &&
t.oldestMs > t.coverageMs
) {
setLimit((l) => l + pageSize);
}
setChildren(ascData);
setHasMoreOlder(saturated);
setIsLoading(false);
setIsLoadingOlder(false);
},
onError: (err) => {
console.warn(err);
setError(err);
setIsLoading(false);
setIsLoadingOlder(false);
},
});
return () => unsubscribe();
}, [collectionPath, limit, pageSize]);
const loadOlder = useCallback(() => {
if (!hasMoreOlder || isLoadingOlder) return;
setIsLoadingOlder(true);
setLimit((l) => l + pageSize);
}, [hasMoreOlder, isLoadingOlder, pageSize]);
if (!path) {
return {
children: [],
isLoading: false,
error: null,
hasMoreOlder: false,
loadOlder: () => {},
isLoadingOlder: false,
};
}
return {
children,
isLoading,
error,
hasMoreOlder,
loadOlder,
isLoadingOlder,
};
}