feat: real-time board updates via WebSocket broadcast + polling

- Add broadcast-service.ts that sends task:changed events to all
  connected WebSocket clients on any task mutation
- Wire broadcasts into task routes (create, update, delete, archive,
  restore, reorder)
- Add useTaskSync hook — connects to WS, listens for task:changed,
  invalidates React Query cache for instant board updates
- Add polling fallback: refetchInterval=10s, staleTime=5s on useTasks
- Auto-reconnect WebSocket with 3s backoff on disconnect

Board now updates within seconds of any external change (API calls
from Clawdbot agents, CLI, etc.) without manual page refresh.

Fixes: task_20260128_Zk4oV0
This commit is contained in:
Brad Groux 2026-01-28 01:20:56 -06:00
parent f7f1487b12
commit 785d60cbcd
8 changed files with 199 additions and 56 deletions

View file

@ -330,3 +330,6 @@
{"type":"task.status_changed","taskId":"task_20260128_tDh8SO","project":"veritas-kanban","status":"in-progress","previousStatus":"todo","id":"evt_oTg72ClmUegb","timestamp":"2026-01-28T07:11:17.460Z"}
{"type":"task.restored","taskId":"task_20260128_cllk71","project":"veritas-kanban","status":"done","id":"evt_4fcYqxmjTMxb","timestamp":"2026-01-28T07:12:02.508Z"}
{"type":"task.archived","taskId":"task_20260128_cllk71","project":"veritas-kanban","status":"done","id":"evt_BiEHiWRStoBu","timestamp":"2026-01-28T07:12:05.263Z"}
{"type":"task.status_changed","taskId":"task_20260128_tDh8SO","project":"veritas-kanban","status":"done","previousStatus":"in-progress","id":"evt_aL7-lVvh-icH","timestamp":"2026-01-28T07:15:45.862Z"}
{"type":"task.status_changed","taskId":"task_20260128_Zk4oV0","project":"veritas-kanban","status":"in-progress","previousStatus":"todo","id":"evt_moXq5uwb9F_S","timestamp":"2026-01-28T07:16:02.124Z"}
{"type":"task.created","taskId":"task_20260128_MnEoHV","status":"todo","id":"evt_w_L_RFmqhIOl","timestamp":"2026-01-28T07:20:14.943Z"}

View file

@ -1,4 +1,55 @@
[
{
"id": "activity_1769584844268_0af9y7c7k",
"type": "task_deleted",
"taskId": "task_20260128_MnEoHV",
"taskTitle": "TEST: Real-time sync verification",
"timestamp": "2026-01-28T07:20:44.269Z"
},
{
"id": "activity_1769584814944_rj5pnt66t",
"type": "task_created",
"taskId": "task_20260128_MnEoHV",
"taskTitle": "TEST: Real-time sync verification",
"details": {
"type": "code",
"priority": "low"
},
"timestamp": "2026-01-28T07:20:14.944Z"
},
{
"id": "activity_1769584562124_j4txg8me8",
"type": "status_changed",
"taskId": "task_20260128_Zk4oV0",
"taskTitle": "BUG: Kanban requires refresh to see progress",
"details": {
"from": "todo",
"status": "in-progress"
},
"timestamp": "2026-01-28T07:16:02.124Z"
},
{
"id": "activity_1769584550594_xx8avtepj",
"type": "comment_added",
"taskId": "task_20260128_tDh8SO",
"taskTitle": "BUG: You can't view or restore archived tasks",
"details": {
"author": "Veritas",
"preview": "Added read-only view for archived tasks (click to ..."
},
"timestamp": "2026-01-28T07:15:50.594Z"
},
{
"id": "activity_1769584545862_qjr4r5px4",
"type": "status_changed",
"taskId": "task_20260128_tDh8SO",
"taskTitle": "BUG: You can't view or restore archived tasks",
"details": {
"from": "in-progress",
"status": "done"
},
"timestamp": "2026-01-28T07:15:45.862Z"
},
{
"id": "activity_1769584325263_iv5d040ns",
"type": "task_archived",
@ -1007,61 +1058,5 @@
"project": "veritas-kanban"
},
"timestamp": "2026-01-28T06:52:50.296Z"
},
{
"id": "activity_1769583170269_4se7u3s82",
"type": "task_created",
"taskId": "task_20260128_Sujk2d",
"taskTitle": "US-401: Diff viewer component",
"details": {
"type": "code",
"priority": "high",
"project": "veritas-kanban"
},
"timestamp": "2026-01-28T06:52:50.269Z"
},
{
"id": "activity_1769583158455_7r6vxey60",
"type": "comment_added",
"taskId": "task_20260128_qImfEu",
"taskTitle": "US-307: Agent completion handling",
"details": {
"author": "Veritas",
"preview": "Historical backfill — completed as part of Sprint ..."
},
"timestamp": "2026-01-28T06:52:38.455Z"
},
{
"id": "activity_1769583158445_7q2a7h0tm",
"type": "status_changed",
"taskId": "task_20260128_qImfEu",
"taskTitle": "US-307: Agent completion handling",
"details": {
"from": "todo",
"status": "done"
},
"timestamp": "2026-01-28T06:52:38.445Z"
},
{
"id": "activity_1769583158435_9sbnapa4g",
"type": "comment_added",
"taskId": "task_20260128_8De44X",
"taskTitle": "US-306: Stop agent",
"details": {
"author": "Veritas",
"preview": "Historical backfill — completed as part of Sprint ..."
},
"timestamp": "2026-01-28T06:52:38.435Z"
},
{
"id": "activity_1769583158426_jcsu1g66m",
"type": "status_changed",
"taskId": "task_20260128_8De44X",
"taskTitle": "US-306: Stop agent",
"details": {
"from": "todo",
"status": "done"
},
"timestamp": "2026-01-28T06:52:38.426Z"
}
]

View file

@ -22,6 +22,7 @@ import metricsRoutes from './routes/metrics.js';
import tracesRoutes from './routes/traces.js';
import attachmentRoutes from './routes/attachments.js';
import { getTelemetryService } from './services/telemetry-service.js';
import { initBroadcast } from './services/broadcast-service.js';
import type { AgentOutput } from './services/agent-service.js';
const app = express();
@ -68,6 +69,9 @@ const server = createServer(app);
// WebSocket server for real-time updates
const wss = new WebSocketServer({ server, path: '/ws' });
// Initialize broadcast service for task change notifications
initBroadcast(wss);
// Track subscriptions: taskId -> Set of WebSocket clients
const agentSubscriptions = new Map<string, Set<WebSocket>>();

View file

@ -4,6 +4,7 @@ import { TaskService } from '../services/task-service.js';
import { WorktreeService } from '../services/worktree-service.js';
import { activityService } from '../services/activity-service.js';
import type { CreateTaskInput, UpdateTaskInput } from '@veritas-kanban/shared';
import { broadcastTaskChange } from '../services/broadcast-service.js';
const router: RouterType = Router();
const taskService = new TaskService();
@ -139,6 +140,7 @@ router.post('/reorder', async (req, res) => {
return res.status(400).json({ error: 'orderedIds must be a non-empty array of task IDs' });
}
const updated = await taskService.reorderTasks(orderedIds);
broadcastTaskChange('reordered');
res.json({ updated: updated.length });
} catch (error) {
console.error('Error reordering tasks:', error);
@ -203,6 +205,7 @@ router.post('/', async (req, res) => {
try {
const input = createTaskSchema.parse(req.body) as CreateTaskInput;
const task = await taskService.createTask(input);
broadcastTaskChange('created', task.id);
// Log activity
await activityService.logActivity('task_created', task.id, task.title, {
@ -248,6 +251,7 @@ router.patch('/:id', async (req, res) => {
if (!task) {
return res.status(404).json({ error: 'Task not found' });
}
broadcastTaskChange('updated', task.id);
// Log activity for status changes
if (input.status && oldTask.status !== input.status) {
@ -277,6 +281,7 @@ router.delete('/:id', async (req, res) => {
if (!success) {
return res.status(404).json({ error: 'Task not found' });
}
broadcastTaskChange('deleted', req.params.id);
// Log activity
if (task) {
@ -298,6 +303,7 @@ router.post('/:id/archive', async (req, res) => {
if (!success) {
return res.status(404).json({ error: 'Task not found' });
}
broadcastTaskChange('archived', req.params.id);
// Log activity
if (task) {
@ -349,6 +355,7 @@ router.post('/:id/restore', async (req, res) => {
if (!task) {
return res.status(404).json({ error: 'Archived task not found' });
}
broadcastTaskChange('restored', task.id);
// Log activity
await activityService.logActivity('status_changed', task.id, task.title, {

View file

@ -0,0 +1,43 @@
import type { WebSocketServer, WebSocket } from 'ws';
/**
* Simple broadcast service that sends task change events to all connected WebSocket clients.
* Initialized with the WebSocketServer instance from index.ts.
*/
let wssRef: WebSocketServer | null = null;
export function initBroadcast(wss: WebSocketServer): void {
wssRef = wss;
}
export type TaskChangeType = 'created' | 'updated' | 'deleted' | 'archived' | 'restored' | 'reordered';
export interface TaskChangeEvent {
type: 'task:changed';
changeType: TaskChangeType;
taskId?: string;
timestamp: string;
}
/**
* Broadcast a task change to all connected WebSocket clients.
* Clients can listen for 'task:changed' messages and invalidate their query caches.
*/
export function broadcastTaskChange(changeType: TaskChangeType, taskId?: string): void {
if (!wssRef) return;
const message: TaskChangeEvent = {
type: 'task:changed',
changeType,
taskId,
timestamp: new Date().toISOString(),
};
const payload = JSON.stringify(message);
wssRef.clients.forEach((client: WebSocket) => {
if (client.readyState === 1) { // WebSocket.OPEN = 1
client.send(payload);
}
});
}

View file

@ -4,8 +4,12 @@ import { Toaster } from './components/ui/toaster';
import { KeyboardProvider } from './hooks/useKeyboard';
import { KeyboardShortcutsDialog } from './components/layout/KeyboardShortcutsDialog';
import { BulkActionsProvider } from './hooks/useBulkActions';
import { useTaskSync } from './hooks/useTaskSync';
function App() {
// Connect to WebSocket for real-time task updates
useTaskSync();
return (
<KeyboardProvider>
<BulkActionsProvider>

View file

@ -0,0 +1,85 @@
import { useEffect, useRef } from 'react';
import { useQueryClient } from '@tanstack/react-query';
/**
* Connects to the Veritas Kanban WebSocket server and listens for
* task:changed events. When received, invalidates the React Query
* task cache so the board updates in real-time.
*/
export function useTaskSync() {
const queryClient = useQueryClient();
const wsRef = useRef<WebSocket | null>(null);
const reconnectTimeoutRef = useRef<number | null>(null);
useEffect(() => {
let mounted = true;
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();
};
}
connect();
return () => {
mounted = false;
if (reconnectTimeoutRef.current) {
clearTimeout(reconnectTimeoutRef.current);
}
if (wsRef.current) {
wsRef.current.close();
}
};
}, [queryClient]);
}

View file

@ -6,6 +6,8 @@ export function useTasks() {
return useQuery({
queryKey: ['tasks'],
queryFn: api.tasks.list,
refetchInterval: 10000, // Poll every 10s as fallback (WebSocket handles instant updates)
staleTime: 5000, // Consider data stale after 5s
});
}