Files
fusion/packages/dashboard/app/hooks/useTasks.ts
gsxdsm 74c0f8fdc2 fix(FN-732): fix dashboard real-time updates and SSE pipeline
- Fix SSE event relay to properly broadcast task store events to dashboard clients
- Use named heartbeat events instead of SSE comments for reliable keep-alive
- Add missing event emission in core task store for state changes
- Add comprehensive tests for SSE pipeline, event emission, and UI hooks
- Remove broken useTerminal hook and AgentLogViewer tests, fix flaky test suites
2026-04-02 19:34:35 -07:00

337 lines
11 KiB
TypeScript

import { useState, useEffect, useCallback, useRef } from "react";
import type { Task, Column, TaskCreateInput, MergeResult } from "@fusion/core";
import * as api from "../api";
const RECONNECT_DELAY_MS = 3000;
/** If no SSE message (including heartbeat events) arrives within this window, force reconnect. */
const HEARTBEAT_TIMEOUT_MS = 45_000;
function normalizeTask(task: Task): Task {
return {
...task,
dependencies: Array.isArray(task.dependencies) ? task.dependencies : [],
steps: Array.isArray(task.steps) ? task.steps : [],
log: Array.isArray((task as Task & { log?: unknown }).log)
? (task as Task & { log?: Task["log"] }).log!
: [],
};
}
/**
* Compare two ISO timestamp strings.
* Returns positive if a is newer than b, negative if b is newer, 0 if equal.
*/
function compareTimestamps(a: string | undefined, b: string | undefined): number {
if (!a && !b) return 0;
if (!a) return -1; // b is newer if a has no timestamp
if (!b) return 1; // a is newer if b has no timestamp
return a.localeCompare(b);
}
export interface UseTasksOptions {
/**
* When provided, fetches tasks only for this project.
* Note: SSE updates are not filtered by project in current implementation.
*/
projectId?: string;
}
export function useTasks(options?: UseTasksOptions) {
const projectId = options?.projectId;
const [tasks, setTasks] = useState<Task[]>([]);
const [connectionNonce, setConnectionNonce] = useState(0);
const tasksRef = useRef(tasks);
const fetchVersionRef = useRef(0);
const lastVisibilityRefreshRef = useRef<number>(0);
tasksRef.current = tasks;
const VISIBILITY_REFRESH_DEBOUNCE_MS = 1000;
const refreshTasks = useCallback(async (options?: { clearOnError?: boolean }) => {
const requestVersion = ++fetchVersionRef.current;
try {
const fetchedTasks = await api.fetchTasks(undefined, undefined, projectId);
if (fetchVersionRef.current !== requestVersion) {
return;
}
setTasks(fetchedTasks.map(normalizeTask));
} catch {
if (fetchVersionRef.current !== requestVersion) {
return;
}
if (options?.clearOnError) {
setTasks([]);
return;
}
setTasks((current) => current);
}
}, [projectId]);
// Fetch initial tasks and recover when the tab becomes visible again.
useEffect(() => {
void refreshTasks({ clearOnError: true });
const handleVisibilityChange = () => {
if (document.visibilityState !== "visible") {
return;
}
const now = Date.now();
const timeSinceLastRefresh = now - lastVisibilityRefreshRef.current;
if (timeSinceLastRefresh < VISIBILITY_REFRESH_DEBOUNCE_MS) {
return;
}
lastVisibilityRefreshRef.current = now;
void refreshTasks();
};
document.addEventListener("visibilitychange", handleVisibilityChange);
return () => {
document.removeEventListener("visibilitychange", handleVisibilityChange);
};
}, [refreshTasks]);
// SSE live updates
// Note: In multi-project mode, SSE receives all task events.
// Tasks are filtered by ID match, so cross-project updates won't affect
// the local state since task IDs are unique and we only fetch from one project.
useEffect(() => {
let closedByCleanup = false;
let reconnectTimer: ReturnType<typeof setTimeout> | null = null;
let heartbeatTimer: ReturnType<typeof setTimeout> | null = null;
if (connectionNonce > 0) {
void refreshTasks();
}
const query = projectId ? `?projectId=${encodeURIComponent(projectId)}` : "";
const es = new EventSource(`/api/events${query}`);
/** Reset the heartbeat watchdog. Called on every incoming SSE message. */
const resetHeartbeat = () => {
if (heartbeatTimer) clearTimeout(heartbeatTimer);
heartbeatTimer = setTimeout(() => {
// No message received within the timeout — connection is likely dead.
if (!closedByCleanup) {
handleError();
}
}, HEARTBEAT_TIMEOUT_MS);
};
// Start the watchdog immediately — if the connection never opens we still want to time out.
resetHeartbeat();
const handleCreated = (e: MessageEvent) => {
resetHeartbeat();
const task = normalizeTask(JSON.parse(e.data) as Task);
// In project mode, only add if this task belongs to our project
// Since we can't determine project from event, we add and let subsequent
// fetches correct the state, or filter by checking if task exists in our set
setTasks((prev) => {
// Avoid duplicates
if (prev.some((t) => t.id === task.id)) return prev;
return [...prev, task];
});
};
const handleMoved = (e: MessageEvent) => {
resetHeartbeat();
const { task, to }: { task: Task; from: Column; to: Column } = JSON.parse(e.data);
const normalizedTask = normalizeTask(task);
setTasks((prev) =>
prev.map((t) =>
t.id === normalizedTask.id ? { ...normalizedTask, column: to } : t
)
);
};
const handleUpdated = (e: MessageEvent) => {
resetHeartbeat();
const incoming = normalizeTask(JSON.parse(e.data) as Task);
setTasks((prev) =>
prev.map((t) => {
if (t.id !== incoming.id) return t;
const updatedAtCompare = compareTimestamps(incoming.updatedAt, t.updatedAt);
if (updatedAtCompare < 0) {
return t;
}
if (t.column === incoming.column) {
return incoming;
}
const columnTimestampCompare = compareTimestamps(t.columnMovedAt, incoming.columnMovedAt);
if (t.columnMovedAt && !incoming.columnMovedAt) {
return { ...incoming, column: t.column, columnMovedAt: t.columnMovedAt };
}
if (columnTimestampCompare > 0) {
return { ...incoming, column: t.column, columnMovedAt: t.columnMovedAt };
}
return incoming;
})
);
};
const handleDeleted = (e: MessageEvent) => {
resetHeartbeat();
const task = normalizeTask(JSON.parse(e.data) as Task);
setTasks((prev) => prev.filter((t) => t.id !== task.id));
};
const handleMerged = (e: MessageEvent) => {
resetHeartbeat();
const { task }: { task: Task } = JSON.parse(e.data);
const normalizedTask = normalizeTask(task);
setTasks((prev) =>
prev.map((t) =>
t.id === normalizedTask.id ? { ...normalizedTask, column: "done" as Column } : t
)
);
};
const cleanup = () => {
if (reconnectTimer) {
clearTimeout(reconnectTimer);
reconnectTimer = null;
}
if (heartbeatTimer) {
clearTimeout(heartbeatTimer);
heartbeatTimer = null;
}
es.removeEventListener("heartbeat", handleHeartbeat);
es.removeEventListener("task:created", handleCreated);
es.removeEventListener("task:moved", handleMoved);
es.removeEventListener("task:updated", handleUpdated);
es.removeEventListener("task:deleted", handleDeleted);
es.removeEventListener("task:merged", handleMerged);
es.removeEventListener("error", handleError);
es.close();
};
const handleError = () => {
if (closedByCleanup) return;
cleanup();
reconnectTimer = setTimeout(() => {
setConnectionNonce((current) => current + 1);
}, RECONNECT_DELAY_MS);
};
/** Server heartbeat (named event, not comment) — just resets the watchdog. */
const handleHeartbeat = () => { resetHeartbeat(); };
es.addEventListener("open", () => resetHeartbeat());
es.addEventListener("heartbeat", handleHeartbeat);
es.addEventListener("task:created", handleCreated);
es.addEventListener("task:moved", handleMoved);
es.addEventListener("task:updated", handleUpdated);
es.addEventListener("task:deleted", handleDeleted);
es.addEventListener("task:merged", handleMerged);
es.addEventListener("error", handleError);
return () => {
closedByCleanup = true;
cleanup();
};
}, [connectionNonce, projectId, refreshTasks]);
const createTask = useCallback(async (input: TaskCreateInput): Promise<Task> => {
const task = normalizeTask(await api.createTask(input, projectId));
setTasks((prev) => {
if (prev.some((t) => t.id === task.id)) return prev;
return [...prev, task];
});
return task;
}, [projectId]);
const moveTask = useCallback(async (id: string, column: Column): Promise<Task> => {
return normalizeTask(await api.moveTask(id, column, projectId));
}, [projectId]);
const deleteTask = useCallback(async (id: string): Promise<Task> => {
return normalizeTask(await api.deleteTask(id, projectId));
}, [projectId]);
const mergeTask = useCallback(async (id: string): Promise<MergeResult> => {
return api.mergeTask(id, projectId);
}, [projectId]);
const retryTask = useCallback(async (id: string): Promise<Task> => {
return normalizeTask(await api.retryTask(id, projectId));
}, [projectId]);
const duplicateTask = useCallback(async (id: string): Promise<Task> => {
const task = normalizeTask(await api.duplicateTask(id, projectId));
setTasks((prev) => {
if (prev.some((t) => t.id === task.id)) return prev;
return [...prev, task];
});
return task;
}, [projectId]);
const updateTask = useCallback(async (
id: string,
updates: { title?: string; description?: string; dependencies?: string[] }
): Promise<Task> => {
const previousTask = tasksRef.current.find((t) => t.id === id);
const optimisticTask = previousTask
? { ...previousTask, ...updates, updatedAt: new Date().toISOString() }
: undefined;
if (optimisticTask) {
setTasks((prev) =>
prev.map((t) => (t.id === id ? optimisticTask : t))
);
}
try {
const updatedTask = normalizeTask(await api.updateTask(id, updates, projectId));
setTasks((prev) =>
prev.map((t) => (t.id === id ? updatedTask : t))
);
return updatedTask;
} catch (err) {
if (previousTask) {
setTasks((prev) =>
prev.map((t) => (t.id === id ? previousTask : t))
);
}
throw err;
}
}, [projectId]);
const archiveTask = useCallback(async (id: string): Promise<Task> => {
const task = normalizeTask(await api.archiveTask(id, projectId));
setTasks((prev) =>
prev.map((t) => (t.id === id ? task : t))
);
return task;
}, [projectId]);
const unarchiveTask = useCallback(async (id: string): Promise<Task> => {
const task = normalizeTask(await api.unarchiveTask(id, projectId));
setTasks((prev) =>
prev.map((t) => (t.id === id ? task : t))
);
return task;
}, [projectId]);
const archiveAllDone = useCallback(async (): Promise<Task[]> => {
const archived = await api.archiveAllDone(projectId);
const normalized = archived.map(normalizeTask);
setTasks((prev) =>
prev.map((t) => {
const updated = normalized.find((archived) => archived.id === t.id);
return updated || t;
})
);
return normalized;
}, [projectId]);
return { tasks, createTask, moveTask, deleteTask, mergeTask, retryTask, duplicateTask, updateTask, archiveTask, unarchiveTask, archiveAllDone };
}