diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 14401e35..4f2decfd 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -96,7 +96,7 @@ jobs: - name: Run Backend Tests run: | ls -la src/generated/prisma - npx vitest run --coverage --reporter=basic + npx vitest run --coverage --reporter=basic --no-file-parallelism working-directory: backend env: DATABASE_URL: postgresql://postgres:password@127.0.0.1:5432/flowfi_test diff --git a/.github/workflows/pr-test-gate.yml b/.github/workflows/pr-test-gate.yml index b2471627..843cb9fe 100644 --- a/.github/workflows/pr-test-gate.yml +++ b/.github/workflows/pr-test-gate.yml @@ -65,7 +65,7 @@ jobs: DATABASE_URL: postgresql://postgres:password@127.0.0.1:5432/flowfi_test - name: Run backend tests - run: npm test + run: npm test -- --no-file-parallelism working-directory: backend env: DATABASE_URL: postgresql://postgres:password@127.0.0.1:5432/flowfi_test diff --git a/backend/tests/integration/stream-lifecycle.test.ts b/backend/tests/integration/stream-lifecycle.test.ts index a374dc42..02625c74 100644 --- a/backend/tests/integration/stream-lifecycle.test.ts +++ b/backend/tests/integration/stream-lifecycle.test.ts @@ -88,7 +88,7 @@ function createStreamCreatedEvent( ]; const overrideEntries: [string, xdr.ScVal][] = Object.entries(overrides).map( - ([k, v]) => [k, nativeToScVal(v)], + ([k, v]) => [k, typeof v === 'bigint' ? scvI128(v) : nativeToScVal(v)], ); return { diff --git a/backend/tests/stream.validator.test.ts b/backend/tests/stream.validator.test.ts index 825e13a8..42ae7dbb 100644 --- a/backend/tests/stream.validator.test.ts +++ b/backend/tests/stream.validator.test.ts @@ -61,6 +61,8 @@ describe('Stream Validator', () => { }; const result = createStreamSchema.safeParse(data); expect(result.success).toBe(false); + expect(result.error).toBeDefined(); + expect(result.error!.issues.length).toBeGreaterThan(0); expect(result.error!.issues[0]!.message).toBe('Rate exceeds maximum allowed value'); }); diff --git a/frontend/src/__tests__/useStreamEvents.test.tsx b/frontend/src/__tests__/useStreamEvents.test.tsx index 9643acb2..06a6dcc8 100644 --- a/frontend/src/__tests__/useStreamEvents.test.tsx +++ b/frontend/src/__tests__/useStreamEvents.test.tsx @@ -9,6 +9,7 @@ type ErrorHandler = () => void; class MockEventSource { static instance: MockEventSource | null = null; + static instanceCount = 0; url: string; onopen: (() => void) | null = null; @@ -21,6 +22,7 @@ class MockEventSource { constructor(url: string) { this.url = url; MockEventSource.instance = this; + MockEventSource.instanceCount += 1; } addEventListener(type: string, handler: EventHandler) { @@ -66,6 +68,7 @@ class MockEventSource { describe('useStreamEvents', () => { beforeEach(() => { MockEventSource.instance = null; + MockEventSource.instanceCount = 0; vi.useFakeTimers(); }); @@ -280,4 +283,53 @@ describe('useStreamEvents', () => { expect(result.current.events).toHaveLength(types.length); }); + + it('creates only one EventSource across multiple re-renders and incoming events', () => { + const { result, rerender } = renderHook( + (opts: { streamIds: string[] } = { streamIds: ['1'] }) => + useStreamEvents({ ...opts, autoReconnect: false }), + ); + + const firstInstance = MockEventSource.instance; + + act(() => { firstInstance?.open(); }); + + // Simulate multiple re-renders with the same subscription (inline array) + rerender({ streamIds: ['1'] }); + rerender({ streamIds: ['1'] }); + rerender({ streamIds: ['1'] }); + + // Simulate incoming events causing re-renders of the consumer + act(() => { + MockEventSource.instance?.emit('stream.created', { i: 1 }); + MockEventSource.instance?.emit('stream.created', { i: 2 }); + MockEventSource.instance?.emit('stream.created', { i: 3 }); + }); + + expect(result.current.events).toHaveLength(3); + + // Re-render again after events + rerender({ streamIds: ['1'] }); + rerender({ streamIds: ['1'] }); + + expect(MockEventSource.instanceCount).toBe(1); + expect(MockEventSource.instance).toBe(firstInstance); + }); + + it('stops reconnecting after reaching the cap', () => { + renderHook(() => + useStreamEvents({ streamIds: ['1'], autoReconnect: true, maxRetryDelay: 1000 }), + ); + + // Trigger errors repeatedly to consume reconnect attempts. + // The reconnect delay stays at 1000ms (capped by maxRetryDelay). + for (let i = 0; i < 25; i++) { + act(() => { MockEventSource.instance?.triggerError(); }); + act(() => { vi.advanceTimersByTime(2000); }); + } + + // 1 initial + 20 reconnect attempts = 21 instances max. + // After the 20th reconnect attempt, no more timers should fire. + expect(MockEventSource.instanceCount).toBeLessThanOrEqual(21); + }); }); diff --git a/frontend/src/hooks/useStreamEvents.ts b/frontend/src/hooks/useStreamEvents.ts index 5cdfc2ce..04f7ff78 100644 --- a/frontend/src/hooks/useStreamEvents.ts +++ b/frontend/src/hooks/useStreamEvents.ts @@ -1,4 +1,4 @@ -import { useEffect, useState, useCallback, useRef } from 'react'; +import { useEffect, useState, useCallback, useRef, useMemo } from 'react'; interface StreamEvent { type: 'created' | 'topped_up' | 'withdrawn' | 'cancelled' | 'completed' | 'paused' | 'resumed'; @@ -23,12 +23,13 @@ interface UseStreamEventsReturn { clearEvents: () => void; } +const MAX_RECONNECT_ATTEMPTS = 20; + export function useStreamEvents( options: UseStreamEventsOptions = {} ): UseStreamEventsReturn { const { - streamIds = [], - // userPublicKeys = [], + streamIds: rawStreamIds = [], subscribeToAll = false, autoReconnect = true, maxRetryDelay = 30000, @@ -43,32 +44,43 @@ export function useStreamEvents( const eventSourceRef = useRef(null); const retryDelayRef = useRef(1000); const reconnectTimeoutRef = useRef | null>(null); + const reconnectAttemptsRef = useRef(0); const connectRef = useRef<() => void>(() => undefined); + const subscriptionKey = useMemo(() => { + const streams = [...rawStreamIds].sort().join(','); + return `${subscribeToAll ? 'all' : streams}|${jwtToken || ''}`; + }, [rawStreamIds, subscribeToAll, jwtToken]); + const buildUrl = useCallback(() => { const params = new URLSearchParams(); if (subscribeToAll) { params.append('all', 'true'); } else { - streamIds.forEach(id => params.append('streams', id)); + rawStreamIds.forEach(id => params.append('streams', id)); } - // Add JWT token to query string for authentication - // (EventSource doesn't support custom headers in browser) if (jwtToken) { params.append('token', jwtToken); } const baseUrl = process.env.NEXT_PUBLIC_API_URL || 'http://localhost:3001'; return `${baseUrl}/v1/events/subscribe?${params}`; - }, [streamIds, subscribeToAll, jwtToken]); + // subscriptionKey captures all subscription parameters as a stable string + // eslint-disable-next-line react-hooks/exhaustive-deps + }, [subscriptionKey]); const clearEvents = useCallback(() => { setEvents([]); }, []); const connect = useCallback(() => { + if (reconnectTimeoutRef.current !== null) { + clearTimeout(reconnectTimeoutRef.current); + reconnectTimeoutRef.current = null; + } + const url = buildUrl(); const eventSource = new EventSource(url); eventSourceRef.current = eventSource; @@ -77,16 +89,16 @@ export function useStreamEvents( setConnected(true); setReconnecting(false); setError(null); - retryDelayRef.current = 1000; // Reset retry delay + retryDelayRef.current = 1000; + reconnectAttemptsRef.current = 0; }; - const handleEvent = (type: StreamEvent['type']) => (e: MessageEvent) => { try { const data = JSON.parse(e.data); setEvents((prev: StreamEvent[]) => [ { type, data, timestamp: Date.now() }, - ...prev.slice(0, 99), // Keep last 100 events + ...prev.slice(0, 99), ]); } catch { // Silently ignore malformed event messages @@ -132,8 +144,9 @@ export function useStreamEvents( eventSourceRef.current.close(); eventSourceRef.current = null; } - if (reconnectTimeoutRef.current) { + if (reconnectTimeoutRef.current !== null) { clearTimeout(reconnectTimeoutRef.current); + reconnectTimeoutRef.current = null; } }; }, [connect]); diff --git a/frontend/src/lib/logger.ts b/frontend/src/lib/logger.ts index 496bfae5..0d1ad792 100644 --- a/frontend/src/lib/logger.ts +++ b/frontend/src/lib/logger.ts @@ -2,16 +2,16 @@ const isDev = process.env.NODE_ENV !== "production"; export const logger = { debug: (...args: unknown[]) => { - if (isDev) console.debug(...args); // eslint-disable-line no-console + if (isDev) console.debug(...args); }, info: (...args: unknown[]) => { - if (isDev) console.info(...args); // eslint-disable-line no-console + if (isDev) console.info(...args); }, warn: (...args: unknown[]) => { - if (isDev) console.warn(...args); // eslint-disable-line no-console + if (isDev) console.warn(...args); }, // errors always surface, even in production error: (...args: unknown[]) => { - console.error(...args); // eslint-disable-line no-console + console.error(...args); }, };