multica/apps/web/features/realtime/use-realtime-sync.ts
Naiyuan Qing 9127e543d5 feat: add event bus, WS workspace isolation, and global store migration
- Add internal event bus (server/internal/events/) with synchronous
  pub/sub and panic isolation per listener
- Upgrade WebSocket Hub to workspace-scoped rooms with JWT auth
  and membership verification on connect
- Add 10 new WS event types (comment CRUD, inbox read/archive,
  agent create/delete, workspace/member events)
- Refactor all handlers and TaskService to publish events via Bus
  instead of direct Hub.Broadcast calls
- Add WS broadcast listener that routes events to correct workspace
- Frontend: WSClient sends token + workspace_id on connect with
  auto-reconnect refetch
- Frontend: centralized useRealtimeSync hook dispatches all WS
  events to global Zustand stores
- Migrate issues and inbox pages from local useState to global
  useIssueStore/useInboxStore
- Make store addIssue/addItem idempotent to prevent duplicates
- Remove dead packages/hooks/src/use-realtime.ts
- Add feature tracking files for 4 planned features

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
2026-03-25 10:08:27 +08:00

115 lines
3.4 KiB
TypeScript

"use client";
import { useEffect } from "react";
import type { WSClient } from "@multica/sdk";
import { useIssueStore, useInboxStore, useAgentStore } from "@multica/store";
import { useWorkspaceStore } from "@/features/workspace";
import { api } from "@/shared/api";
import type {
IssueCreatedPayload,
IssueUpdatedPayload,
IssueDeletedPayload,
AgentStatusPayload,
AgentCreatedPayload,
InboxNewPayload,
InboxReadPayload,
InboxArchivedPayload,
} from "@multica/types";
/**
* Centralized WS → store sync. Called once from WSProvider.
* Subscribes to all global WS events and dispatches to Zustand stores.
* Comment events are NOT handled here — they stay per-page on issue detail.
*/
export function useRealtimeSync(ws: WSClient | null) {
// Issue events → useIssueStore
useEffect(() => {
if (!ws) return;
const unsubs = [
ws.on("issue:created", (p) => {
const { issue } = p as IssueCreatedPayload;
useIssueStore.getState().addIssue(issue);
}),
ws.on("issue:updated", (p) => {
const { issue } = p as IssueUpdatedPayload;
useIssueStore.getState().updateIssue(issue.id, issue);
}),
ws.on("issue:deleted", (p) => {
const { issue_id } = p as IssueDeletedPayload;
useIssueStore.getState().removeIssue(issue_id);
}),
];
return () => unsubs.forEach((u) => u());
}, [ws]);
// Inbox events → useInboxStore
useEffect(() => {
if (!ws) return;
const unsubs = [
ws.on("inbox:new", (p) => {
const { item } = p as InboxNewPayload;
useInboxStore.getState().addItem(item);
}),
ws.on("inbox:read", (p) => {
const { item_id } = p as InboxReadPayload;
useInboxStore.getState().markRead(item_id);
}),
ws.on("inbox:archived", (p) => {
const { item_id } = p as InboxArchivedPayload;
useInboxStore.getState().archive(item_id);
}),
];
return () => unsubs.forEach((u) => u());
}, [ws]);
// Agent events → useAgentStore / workspace refresh
useEffect(() => {
if (!ws) return;
const unsubs = [
ws.on("agent:status", (p) => {
const { agent } = p as AgentStatusPayload;
useAgentStore.getState().updateAgent(agent.id, agent);
}),
ws.on("agent:created", (p) => {
const { agent } = p as AgentCreatedPayload;
const agents = useAgentStore.getState().agents;
if (!agents.find((a) => a.id === agent.id)) {
useAgentStore.getState().setAgents([...agents, agent]);
}
}),
ws.on("agent:deleted", () => {
// Refresh agents list since we don't have removeAgent in store
useWorkspaceStore.getState().refreshAgents();
}),
];
return () => unsubs.forEach((u) => u());
}, [ws]);
// Reconnect → refetch all data to recover missed events
useEffect(() => {
if (!ws) return;
const unsub = ws.onReconnect(async () => {
try {
const [issuesRes, inboxItems, agents] = await Promise.all([
api.listIssues({ limit: 200 }),
api.listInbox(),
api.listAgents(),
]);
useIssueStore.getState().setIssues(issuesRes.issues);
useInboxStore.getState().setItems(inboxItems);
useAgentStore.getState().setAgents(agents);
} catch {
// Silently fail; next reconnect will retry
}
});
return unsub;
}, [ws]);
}