plus-ui/src/chat-core/sse/SseSubscriber.ts
i548450 5c8b6dc014 feat(customerservice): chat-core + chat-ui + workbench 三栏
- chat-core: 类型定义(Message/Conversation/SsePayload) + SseSubscriber
  订阅共享 sseEventBus 'customerservice' event,按 msgType 分发
- chat-core utils: formatRelativeTime / extractUrlsAsLinks + 19 个 vitest 单测
- chat-ui: MessageStream/MessageBubble/MessageInput/MessageStateBadge/
  DateDivider/ConversationDivider 组件 + useSseConversation/useMessageStream
  composables;CSS variable 暴露主题色
- api/customerservice: conversation/message/agent/mediaUser HTTP 客户端
  按 spec 04-api.md endpoint 形态先行,后端 Task 5-13 完成后无需改前端
- views/customerservice/workbench: 三栏布局 + Pinia store 处理 6 种 SSE
  payload + 自动重连后全量对齐(spec 05-sse.md §6)
- vitest.config.ts: 独立 standalone 配置,跳过 dev server 的 unocss 链
2026-06-09 21:45:34 +08:00

74 lines
2.3 KiB
TypeScript

import type { Emitter } from 'mitt';
import { SSE_BUS_EVENT, type SseHandler, type SseMsgType, type SsePayload } from '../types/sse';
/**
* Subscribes to a shared SSE event bus (sse.ts → mitt) and dispatches
* customerservice payloads to msgType-keyed handlers.
*
* Does NOT create its own EventSource. The host bus is provided by
* frontend/src/utils/sseEventBus.getSseEventBus().
*/
export class SseSubscriber {
private readonly handlers: Map<SseMsgType, Set<(data: any) => void>> = new Map();
private readonly reconnectHandlers: Set<() => void> = new Set();
private readonly busListener: (raw: string) => void;
private readonly reconnectListener: () => void;
private disposed = false;
constructor(private readonly bus: Emitter<any>) {
this.busListener = (raw: string) => this.handleRaw(raw);
this.reconnectListener = () => this.reconnectHandlers.forEach((h) => h());
this.bus.on(SSE_BUS_EVENT, this.busListener);
this.bus.on('reconnect', this.reconnectListener);
}
on<K extends SseMsgType>(msgType: K, handler: SseHandler<K>): void {
if (this.disposed) return;
let set = this.handlers.get(msgType);
if (!set) {
set = new Set();
this.handlers.set(msgType, set);
}
set.add(handler as (data: any) => void);
}
off<K extends SseMsgType>(msgType: K, handler: SseHandler<K>): void {
this.handlers.get(msgType)?.delete(handler as (data: any) => void);
}
onReconnect(handler: () => void): void {
if (this.disposed) return;
this.reconnectHandlers.add(handler);
}
dispose(): void {
this.bus.off(SSE_BUS_EVENT, this.busListener);
this.bus.off('reconnect', this.reconnectListener);
this.handlers.clear();
this.reconnectHandlers.clear();
this.disposed = true;
}
private handleRaw(raw: string): void {
let envelope: SsePayload | null = null;
try {
const parsed = JSON.parse(raw);
if (parsed && typeof parsed === 'object' && typeof parsed.msgType === 'string') {
envelope = parsed as SsePayload;
}
} catch {
return;
}
if (!envelope) return;
const set = this.handlers.get(envelope.msgType);
if (!set) return;
set.forEach((h) => {
try {
h(envelope!.data);
} catch (err) {
console.error('[SseSubscriber] handler threw for', envelope!.msgType, err);
}
});
}
}