RF-12: Extract shared useWebSocket hook

- 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
This commit is contained in:
Brad Groux 2026-01-28 06:33:29 -06:00
parent 58bfae2c25
commit ada2e468a0
6 changed files with 274 additions and 129 deletions

View file

@ -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"}

View file

@ -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",

View file

@ -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';

View file

@ -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<AgentOutput[]>([]);
const [isConnected, setIsConnected] = useState(false);
const [isRunning, setIsRunning] = useState(false);
const wsRef = useRef<WebSocket | null>(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([]);
}, []);

View file

@ -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<WebSocket | null>(null);
const reconnectTimeoutRef = useRef<number | null>(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,
});
}

View file

@ -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<WebSocketMessage | null>(null);
const wsRef = useRef<WebSocket | null>(null);
const reconnectTimeoutRef = useRef<ReturnType<typeof setTimeout> | 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,
};
}