'use client'; import { Pause, Play, RotateCcw } from 'lucide-react'; import { useSearchParams } from 'next/navigation'; import { useCallback, useEffect, useMemo, useRef, useState } from 'react'; import { useEventDrawer } from '@/components/events/event-drawer-context'; import { EventRow } from '@/components/events/event-row'; import { LiveStatus } from '@/components/ui/live'; import { clientApi } from '@/lib/client-api'; import { cn } from '@/lib/cn'; import type { Event } from '@/lib/types'; const MAX_ROWS = 300; function matches(e: Event, f: { event_type?: string; min_importance?: number; country?: string; industry?: string; min_confidence?: number }): boolean { if (f.event_type && e.event_type !== f.event_type) return false; if (f.min_importance !== undefined && (e.importance > 1 ? e.importance / 100 : e.importance) < f.min_importance) return false; if (f.min_confidence !== undefined && e.confidence < f.min_confidence) return false; if (f.country && (e.company.country ?? '').toUpperCase() !== f.country.toUpperCase()) return false; return true; } /** * Live Activity Feed. Seeds with server-rendered events, then subscribes to `/api/v1/live/stream` (SSE, `event: event`). * New rows prepend with a soft highlight; when paused (hover or button) they buffer behind an "N new" pill. * Keyboard j/k moves a selection, Enter opens the evidence drawer. Filters come from the URL (`?event_type=&min_importance=…`) * and are applied client-side to the stream as well. Falls back to polling `/live?since=` when EventSource fails. */ export function LiveFeed({ initial, limit = 60, compact = false, className, showControls = true, stickyPill = true }: { initial: Event[]; limit?: number; compact?: boolean; className?: string; showControls?: boolean; stickyPill?: boolean }) { const sp = useSearchParams(); const filters = useMemo( () => ({ event_type: sp.get('event_type') ?? undefined, min_importance: sp.get('min_importance') ? Number(sp.get('min_importance')) : undefined, min_confidence: sp.get('min_confidence') ? Number(sp.get('min_confidence')) : undefined, country: sp.get('country') ?? undefined, industry: sp.get('industry') ?? undefined, }), [sp], ); const [events, setEvents] = useState(() => initial.filter((e) => matches(e, filters)).slice(0, limit)); const [buffer, setBuffer] = useState([]); const [fresh, setFresh] = useState>(new Set()); const [paused, setPaused] = useState(false); const [hovering, setHovering] = useState(false); const [connected, setConnected] = useState(false); const [updatedAt, setUpdatedAt] = useState(null); const [selected, setSelected] = useState(-1); const seen = useRef>(new Set(initial.map((e) => e.id))); const latest = useRef(initial[0]?.detected_at ?? null); const { open } = useEventDrawer(); const listRef = useRef(null); const holding = paused || hovering; const holdingRef = useRef(holding); holdingRef.current = holding; const filtersRef = useRef(filters); filtersRef.current = filters; // re-seed when filters change (server re-renders `initial` on navigation) useEffect(() => { setEvents(initial.filter((e) => matches(e, filters)).slice(0, limit)); seen.current = new Set(initial.map((e) => e.id)); setBuffer([]); }, [initial, filters, limit]); const ingest = useCallback( (incoming: Event[]) => { const fresh = incoming.filter((e) => !seen.current.has(e.id) && matches(e, filtersRef.current)); if (!fresh.length) return; for (const e of fresh) seen.current.add(e.id); const newest = fresh.map((e) => e.detected_at).sort().pop(); if (newest && (!latest.current || newest > latest.current)) latest.current = newest; setUpdatedAt(Date.now()); if (holdingRef.current) { setBuffer((b) => [...fresh, ...b].slice(0, MAX_ROWS)); return; } setEvents((prev) => [...fresh, ...prev].slice(0, Math.max(limit, MAX_ROWS))); setFresh((s) => { const n = new Set(s); for (const e of fresh) n.add(e.id); return n; }); setTimeout(() => setFresh((s) => { const n = new Set(s); for (const e of fresh) n.delete(e.id); return n; }), 2000); }, [limit], ); // SSE with polling fallback useEffect(() => { let es: EventSource | null = null; let poll: ReturnType | null = null; let cancelled = false; const qs = new URLSearchParams(); if (latest.current) qs.set('since', latest.current); if (filters.event_type) qs.set('event_type', filters.event_type); if (filters.min_importance !== undefined) qs.set('min_importance', String(filters.min_importance)); const startPolling = () => { if (poll) return; poll = setInterval(async () => { try { const res = await clientApi.live({ since: latest.current ?? undefined, limit: 50, event_type: filters.event_type }); const items = Array.isArray(res) ? res : res.items; if (!cancelled) { setConnected(true); ingest(items ?? []); } } catch { if (!cancelled) setConnected(false); } }, 8000); }; // Watchdog: a proxy that gzips text/event-stream buffers frames indefinitely (the API must send `no-transform`); // if the socket is "open" but silent for 45 s, poll as well so the feed keeps moving. let lastFrame = Date.now(); const watchdog = setInterval(() => { if (Date.now() - lastFrame > 45_000) startPolling(); }, 15_000); if (typeof EventSource !== 'undefined') { try { es = new EventSource(`/api/v1/live/stream${qs.toString() ? `?${qs}` : ''}`); es.addEventListener('open', () => { setConnected(true); setUpdatedAt(Date.now()); }); es.addEventListener('event', (m) => { lastFrame = Date.now(); try { const e = JSON.parse((m as MessageEvent).data) as Event; ingest([e]); } catch { /* malformed frame */ } }); es.addEventListener('heartbeat', () => { lastFrame = Date.now(); setUpdatedAt(Date.now()); }); es.onerror = () => { setConnected(false); // EventSource retries by itself; also start a low-frequency poll so the feed keeps moving behind proxies. startPolling(); }; } catch { startPolling(); } } else startPolling(); return () => { cancelled = true; es?.close(); clearInterval(watchdog); if (poll) clearInterval(poll); }; }, [filters, ingest]); const release = () => { if (!buffer.length) return; setEvents((prev) => [...buffer, ...prev].slice(0, MAX_ROWS)); setFresh((s) => { const n = new Set(s); for (const e of buffer) n.add(e.id); return n; }); const ids = buffer.map((e) => e.id); setTimeout(() => setFresh((s) => { const n = new Set(s); for (const id of ids) n.delete(id); return n; }), 2000); setBuffer([]); }; useEffect(() => { if (!holding && buffer.length) release(); // eslint-disable-next-line react-hooks/exhaustive-deps }, [holding]); // keyboard j/k/Enter useEffect(() => { const onKey = (e: KeyboardEvent) => { const t = e.target as HTMLElement | null; if (t && (t.tagName === 'INPUT' || t.tagName === 'TEXTAREA' || t.tagName === 'SELECT' || t.isContentEditable)) return; if (document.querySelector('[role="dialog"]')) return; if (e.key === 'j' || e.key === 'k') { e.preventDefault(); setSelected((s) => { const n = e.key === 'j' ? Math.min(events.length - 1, s + 1) : Math.max(0, s - 1); listRef.current?.querySelector(`[data-index="${n}"]`)?.scrollIntoView({ block: 'nearest' }); return n; }); } else if (e.key === 'Enter' && selected >= 0 && events[selected]) { e.preventDefault(); open(events[selected] as Event); } }; window.addEventListener('keydown', onKey); return () => window.removeEventListener('keydown', onKey); }, [events, selected, open]); const replay = async () => { const since = new Date(Date.now() - 3600_000).toISOString(); try { const res = await clientApi.live({ since, limit: 100, event_type: filters.event_type }); const items = Array.isArray(res) ? res : res.items; ingest(items ?? []); } catch { /* ignore */ } }; return (
setHovering(true)} onMouseLeave={() => setHovering(false)}> {showControls && (
{holding ? (paused ? 'paused' : 'paused while hovering') : ''}
)} {buffer.length > 0 && (
)} {events.length === 0 ? (

No monitored evidence matches these filters yet. The feed stays connected — new events will appear here.

) : (
    {events.slice(0, compact ? limit : MAX_ROWS).map((e, i) => (
  • ))}
)} {!compact &&

Keyboard: j / k to move · Enter to open evidence · hover pauses the stream.

}
); }