From ada2e468a0da7cf2329f3fe4cbe246dc1211c779 Mon Sep 17 00:00:00 2001 From: Brad Groux Date: Wed, 28 Jan 2026 06:33:29 -0600 Subject: [PATCH] RF-12: Extract shared useWebSocket hook MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Created hooks/useWebSocket.ts with connection lifecycle, reconnection, typed messages - Refactored useTaskSync.ts to use shared hook (80 → 35 lines) - Refactored useAgent.ts to use shared hook (cleaner separation) - Added useWebSocket to barrel export --- .../telemetry/events-2026-01-28.ndjson | 2 + server/.veritas-kanban/activity.json | 33 ++++ web/src/hooks/index.ts | 1 + web/src/hooks/useAgent.ts | 88 +++------ web/src/hooks/useTaskSync.ts | 93 +++------ web/src/hooks/useWebSocket.ts | 186 ++++++++++++++++++ 6 files changed, 274 insertions(+), 129 deletions(-) create mode 100644 web/src/hooks/useWebSocket.ts diff --git a/.veritas-kanban/telemetry/events-2026-01-28.ndjson b/.veritas-kanban/telemetry/events-2026-01-28.ndjson index 5adca4a7..6d05a848 100644 --- a/.veritas-kanban/telemetry/events-2026-01-28.ndjson +++ b/.veritas-kanban/telemetry/events-2026-01-28.ndjson @@ -598,3 +598,5 @@ {"type":"task.status_changed","taskId":"task_20260128_oeprdJ","project":"veritas-kanban","status":"in-progress","previousStatus":"todo","id":"evt_f00PlGFiT8-I","timestamp":"2026-01-28T12:25:50.647Z"} {"type":"task.status_changed","taskId":"task_20260128_oeprdJ","project":"veritas-kanban","status":"done","previousStatus":"in-progress","id":"evt_mUOZJIeRjCgA","timestamp":"2026-01-28T12:29:24.872Z"} {"type":"task.status_changed","taskId":"task_20260128_Qh60Ao","project":"veritas-kanban","status":"in-progress","previousStatus":"todo","id":"evt_IMzc-x6hnSTj","timestamp":"2026-01-28T12:29:48.687Z"} +{"type":"task.status_changed","taskId":"task_20260128_Qh60Ao","project":"veritas-kanban","status":"done","previousStatus":"in-progress","id":"evt_SljUoovW7DVZ","timestamp":"2026-01-28T12:31:58.645Z"} +{"type":"task.status_changed","taskId":"task_20260128_VnO-L8","project":"veritas-kanban","status":"in-progress","previousStatus":"todo","id":"evt_pCeJ0lWaPH9N","timestamp":"2026-01-28T12:32:22.407Z"} diff --git a/server/.veritas-kanban/activity.json b/server/.veritas-kanban/activity.json index 5cfb046b..bb0baf05 100644 --- a/server/.veritas-kanban/activity.json +++ b/server/.veritas-kanban/activity.json @@ -1,4 +1,37 @@ [ + { + "id": "activity_1769603542408_0wffse0lg", + "type": "status_changed", + "taskId": "task_20260128_VnO-L8", + "taskTitle": "RF-12: Extract Shared WebSocket Hook", + "details": { + "from": "todo", + "status": "in-progress" + }, + "timestamp": "2026-01-28T12:32:22.408Z" + }, + { + "id": "activity_1769603526556_yx9ze60cc", + "type": "comment_added", + "taskId": "task_20260128_Qh60Ao", + "taskTitle": "RF-22: TypeScript Strictness & Linting", + "details": { + "author": "Veritas", + "preview": "Added ESLint with TypeScript and React plugins. Al..." + }, + "timestamp": "2026-01-28T12:32:06.556Z" + }, + { + "id": "activity_1769603518646_2mawe9n8u", + "type": "status_changed", + "taskId": "task_20260128_Qh60Ao", + "taskTitle": "RF-22: TypeScript Strictness & Linting", + "details": { + "from": "in-progress", + "status": "done" + }, + "timestamp": "2026-01-28T12:31:58.646Z" + }, { "id": "activity_1769603388688_4hoz7af0m", "type": "status_changed", diff --git a/web/src/hooks/index.ts b/web/src/hooks/index.ts index c6ee7dfe..78dec1ae 100644 --- a/web/src/hooks/index.ts +++ b/web/src/hooks/index.ts @@ -36,4 +36,5 @@ export { } from './useTimeTracking'; // formatDuration from useTimeTracking takes seconds; use useMetrics.formatDuration (takes ms) via barrel export * from './useToast'; +export * from './useWebSocket'; export * from './useWorktree'; diff --git a/web/src/hooks/useAgent.ts b/web/src/hooks/useAgent.ts index 2967ab2a..04c591b9 100644 --- a/web/src/hooks/useAgent.ts +++ b/web/src/hooks/useAgent.ts @@ -1,6 +1,7 @@ -import { useState, useEffect, useCallback, useRef } from 'react'; +import { useState, useCallback, useEffect } from 'react'; import { useQuery, useMutation, useQueryClient } from '@tanstack/react-query'; import { api, AgentOutput } from '@/lib/api'; +import { useWebSocket, type WebSocketMessage } from './useWebSocket'; import type { AgentType } from '@veritas-kanban/shared'; export function useAgentStatus(taskId: string | undefined) { @@ -63,69 +64,40 @@ export function useAgentLog(taskId: string | undefined, attemptId: string | unde // WebSocket hook for real-time agent output export function useAgentStream(taskId: string | undefined) { const [outputs, setOutputs] = useState([]); - const [isConnected, setIsConnected] = useState(false); const [isRunning, setIsRunning] = useState(false); - const wsRef = useRef(null); const queryClient = useQueryClient(); - useEffect(() => { - if (!taskId) { - setOutputs([]); - return; + const handleMessage = useCallback((message: WebSocketMessage) => { + if (message.type === 'subscribed') { + setIsRunning(message.running as boolean); + } else if (message.type === 'agent:output') { + setOutputs(prev => [...prev, { + type: message.outputType as AgentOutput['type'], + content: message.content as string, + timestamp: message.timestamp as string, + }]); + } else if (message.type === 'agent:complete') { + setIsRunning(false); + queryClient.invalidateQueries({ queryKey: ['agent', 'status', taskId] }); + queryClient.invalidateQueries({ queryKey: ['tasks'] }); + } else if (message.type === 'agent:error') { + setIsRunning(false); + queryClient.invalidateQueries({ queryKey: ['agent', 'status', taskId] }); } - - // Connect to WebSocket - const protocol = window.location.protocol === 'https:' ? 'wss:' : 'ws:'; - const wsUrl = `${protocol}//${window.location.hostname}:3001/ws`; - - const ws = new WebSocket(wsUrl); - wsRef.current = ws; - - ws.onopen = () => { - setIsConnected(true); - // Subscribe to task's agent output - ws.send(JSON.stringify({ type: 'subscribe', taskId })); - }; - - ws.onmessage = (event) => { - try { - const data = JSON.parse(event.data); - - if (data.type === 'subscribed') { - setIsRunning(data.running); - } else if (data.type === 'agent:output') { - setOutputs(prev => [...prev, { - type: data.type === 'agent:output' ? data.type : data.type, - content: data.content, - timestamp: data.timestamp, - } as AgentOutput]); - } else if (data.type === 'agent:complete') { - setIsRunning(false); - queryClient.invalidateQueries({ queryKey: ['agent', 'status', taskId] }); - queryClient.invalidateQueries({ queryKey: ['tasks'] }); - } else if (data.type === 'agent:error') { - setIsRunning(false); - queryClient.invalidateQueries({ queryKey: ['agent', 'status', taskId] }); - } - } catch (e) { - console.error('WebSocket message parse error:', e); - } - }; - - ws.onclose = () => { - setIsConnected(false); - }; - - ws.onerror = (error) => { - console.error('WebSocket error:', error); - }; - - return () => { - ws.close(); - wsRef.current = null; - }; }, [taskId, queryClient]); + // Clear outputs when taskId changes + useEffect(() => { + setOutputs([]); + }, [taskId]); + + const { isConnected } = useWebSocket({ + autoConnect: !!taskId, + onOpen: taskId ? { type: 'subscribe', taskId } : undefined, + onMessage: handleMessage, + reconnectDelay: 0, // Don't reconnect for agent streams + }); + const clearOutputs = useCallback(() => { setOutputs([]); }, []); diff --git a/web/src/hooks/useTaskSync.ts b/web/src/hooks/useTaskSync.ts index fe9cc788..d4015859 100644 --- a/web/src/hooks/useTaskSync.ts +++ b/web/src/hooks/useTaskSync.ts @@ -1,5 +1,6 @@ -import { useEffect, useRef } from 'react'; +import { useCallback } from 'react'; import { useQueryClient } from '@tanstack/react-query'; +import { useWebSocket, type WebSocketMessage } from './useWebSocket'; /** * Connects to the Veritas Kanban WebSocket server and listens for @@ -8,78 +9,28 @@ import { useQueryClient } from '@tanstack/react-query'; */ export function useTaskSync() { const queryClient = useQueryClient(); - const wsRef = useRef(null); - const reconnectTimeoutRef = useRef(null); - useEffect(() => { - let mounted = true; + const handleMessage = useCallback((message: WebSocketMessage) => { + if (message.type === 'task:changed') { + // Invalidate task queries to trigger a refetch + queryClient.invalidateQueries({ queryKey: ['tasks'] }); + + // Also invalidate specific task if we know which one + if (message.taskId) { + queryClient.invalidateQueries({ queryKey: ['tasks', message.taskId] }); + } - function connect() { - if (!mounted) return; - - // Determine WebSocket URL from current location - const protocol = window.location.protocol === 'https:' ? 'wss:' : 'ws:'; - // In dev (port 3000/5173), the API is on port 3001; in prod it's same host - const isDev = ['3000', '5173'].includes(window.location.port); - const wsHost = isDev - ? `${protocol}//localhost:3001/ws` - : `${protocol}//${window.location.host}/ws`; - - const ws = new WebSocket(wsHost); - wsRef.current = ws; - - ws.onopen = () => { - // Send a subscribe message for task changes - ws.send(JSON.stringify({ type: 'subscribe:tasks' })); - }; - - ws.onmessage = (event) => { - try { - const message = JSON.parse(event.data); - - if (message.type === 'task:changed') { - // Invalidate task queries to trigger a refetch - queryClient.invalidateQueries({ queryKey: ['tasks'] }); - - // Also invalidate specific task if we know which one - if (message.taskId) { - queryClient.invalidateQueries({ queryKey: ['tasks', message.taskId] }); - } - - // If it's an archive-related change, also invalidate archive queries - if (message.changeType === 'archived' || message.changeType === 'restored') { - queryClient.invalidateQueries({ queryKey: ['tasks', 'archived'] }); - queryClient.invalidateQueries({ queryKey: ['tasks', 'archive-suggestions'] }); - } - } - } catch { - // Ignore parse errors for non-task messages - } - }; - - ws.onclose = () => { - wsRef.current = null; - // Reconnect after 3 seconds - if (mounted) { - reconnectTimeoutRef.current = window.setTimeout(connect, 3000); - } - }; - - ws.onerror = () => { - ws.close(); - }; + // If it's an archive-related change, also invalidate archive queries + if (message.changeType === 'archived' || message.changeType === 'restored') { + queryClient.invalidateQueries({ queryKey: ['tasks', 'archived'] }); + queryClient.invalidateQueries({ queryKey: ['tasks', 'archive-suggestions'] }); + } } - - connect(); - - return () => { - mounted = false; - if (reconnectTimeoutRef.current) { - clearTimeout(reconnectTimeoutRef.current); - } - if (wsRef.current) { - wsRef.current.close(); - } - }; }, [queryClient]); + + useWebSocket({ + onOpen: { type: 'subscribe:tasks' }, + onMessage: handleMessage, + reconnectDelay: 3000, + }); } diff --git a/web/src/hooks/useWebSocket.ts b/web/src/hooks/useWebSocket.ts new file mode 100644 index 00000000..af8f1a5a --- /dev/null +++ b/web/src/hooks/useWebSocket.ts @@ -0,0 +1,186 @@ +import { useEffect, useRef, useState, useCallback } from 'react'; + +export interface WebSocketMessage { + type: string; + [key: string]: unknown; +} + +export interface UseWebSocketOptions { + /** URL to connect to. Defaults to ws(s)://host:3001/ws */ + url?: string; + /** Whether to automatically connect. Default true. */ + autoConnect?: boolean; + /** Message to send on open (subscription). */ + onOpen?: WebSocketMessage; + /** Reconnect delay in ms. 0 to disable. Default 3000. */ + reconnectDelay?: number; + /** Maximum reconnect attempts. 0 for unlimited. Default 0. */ + maxReconnectAttempts?: number; + /** Callback when connection opens. */ + onConnected?: () => void; + /** Callback when connection closes. */ + onDisconnected?: () => void; + /** Message handler. */ + onMessage?: (message: WebSocketMessage) => void; + /** Error handler. */ + onError?: (error: Event) => void; +} + +export interface UseWebSocketReturn { + /** Whether currently connected. */ + isConnected: boolean; + /** Send a message. */ + send: (message: WebSocketMessage) => void; + /** Manually connect. */ + connect: () => void; + /** Manually disconnect. */ + disconnect: () => void; + /** Last received message. */ + lastMessage: WebSocketMessage | null; +} + +function getDefaultWsUrl(): string { + const protocol = window.location.protocol === 'https:' ? 'wss:' : 'ws:'; + const isDev = ['3000', '5173'].includes(window.location.port); + return isDev + ? `${protocol}//localhost:3001/ws` + : `${protocol}//${window.location.host}/ws`; +} + +export function useWebSocket(options: UseWebSocketOptions = {}): UseWebSocketReturn { + const { + url, + autoConnect = true, + onOpen, + reconnectDelay = 3000, + maxReconnectAttempts = 0, + onConnected, + onDisconnected, + onMessage, + onError, + } = options; + + const [isConnected, setIsConnected] = useState(false); + const [lastMessage, setLastMessage] = useState(null); + + const wsRef = useRef(null); + const reconnectTimeoutRef = useRef | null>(null); + const reconnectAttemptsRef = useRef(0); + const mountedRef = useRef(true); + + // Store callbacks in refs to avoid reconnecting on callback changes + const onConnectedRef = useRef(onConnected); + const onDisconnectedRef = useRef(onDisconnected); + const onMessageRef = useRef(onMessage); + const onErrorRef = useRef(onError); + const onOpenRef = useRef(onOpen); + + useEffect(() => { + onConnectedRef.current = onConnected; + onDisconnectedRef.current = onDisconnected; + onMessageRef.current = onMessage; + onErrorRef.current = onError; + onOpenRef.current = onOpen; + }, [onConnected, onDisconnected, onMessage, onError, onOpen]); + + const clearReconnectTimeout = useCallback(() => { + if (reconnectTimeoutRef.current) { + clearTimeout(reconnectTimeoutRef.current); + reconnectTimeoutRef.current = null; + } + }, []); + + const connect = useCallback(() => { + if (!mountedRef.current) return; + if (wsRef.current?.readyState === WebSocket.OPEN) return; + + clearReconnectTimeout(); + + const wsUrl = url || getDefaultWsUrl(); + const ws = new WebSocket(wsUrl); + wsRef.current = ws; + + ws.onopen = () => { + if (!mountedRef.current) return; + setIsConnected(true); + reconnectAttemptsRef.current = 0; + onConnectedRef.current?.(); + + // Send subscription message if provided + if (onOpenRef.current) { + ws.send(JSON.stringify(onOpenRef.current)); + } + }; + + ws.onmessage = (event) => { + if (!mountedRef.current) return; + try { + const message = JSON.parse(event.data) as WebSocketMessage; + setLastMessage(message); + onMessageRef.current?.(message); + } catch (e) { + console.error('WebSocket message parse error:', e); + } + }; + + ws.onclose = () => { + if (!mountedRef.current) return; + setIsConnected(false); + wsRef.current = null; + onDisconnectedRef.current?.(); + + // Reconnect if enabled + if (reconnectDelay > 0) { + const canRetry = maxReconnectAttempts === 0 || + reconnectAttemptsRef.current < maxReconnectAttempts; + + if (canRetry) { + reconnectAttemptsRef.current++; + reconnectTimeoutRef.current = setTimeout(connect, reconnectDelay); + } + } + }; + + ws.onerror = (error) => { + onErrorRef.current?.(error); + ws.close(); + }; + }, [url, reconnectDelay, maxReconnectAttempts, clearReconnectTimeout]); + + const disconnect = useCallback(() => { + clearReconnectTimeout(); + reconnectAttemptsRef.current = maxReconnectAttempts; // Prevent auto-reconnect + wsRef.current?.close(); + wsRef.current = null; + }, [clearReconnectTimeout, maxReconnectAttempts]); + + const send = useCallback((message: WebSocketMessage) => { + if (wsRef.current?.readyState === WebSocket.OPEN) { + wsRef.current.send(JSON.stringify(message)); + } + }, []); + + // Connect on mount if autoConnect + useEffect(() => { + mountedRef.current = true; + + if (autoConnect) { + connect(); + } + + return () => { + mountedRef.current = false; + clearReconnectTimeout(); + wsRef.current?.close(); + wsRef.current = null; + }; + }, [autoConnect, connect, clearReconnectTimeout]); + + return { + isConnected, + send, + connect, + disconnect, + lastMessage, + }; +}