import type { SuperagentAgentDonePayload, SuperagentConversation, SuperagentMessage, SuperagentRealtimeClient, SuperagentRealtimeHandlers, } from '../types'; type SocketListener = (...args: unknown[]) => void; export type SuperagentSocketLike = { connected?: boolean; emit(event: string, ...args: unknown[]): void; on(event: string, listener: SocketListener): void; off?(event: string, listener: SocketListener): void; removeListener?(event: string, listener: SocketListener): void; }; export type SuperagentSocketClientConfig = { socket: SuperagentSocketLike; roomMetadata?: Record | ((conversationId: string) => Record | undefined); }; type RealtimeEnvelope = { room: string; data?: unknown; }; export function createSuperagentSocketClient({ socket, roomMetadata, }: SuperagentSocketClientConfig): SuperagentRealtimeClient { const roomCounts = new Map(); return { subscribeToConversation(conversationId, handlers) { const room = getConversationRoom(conversationId); const metadata = typeof roomMetadata === 'function' ? roomMetadata(conversationId) : roomMetadata; addRoom(room, metadata); const onUpdate: SocketListener = (raw) => { const envelope = parseEnvelope(raw); if (!envelope || envelope.room !== room) return; dispatchUpdate(envelope.data, handlers); }; const onAgentDone: SocketListener = (raw) => { const envelope = parseEnvelope(raw); if (!envelope || envelope.room !== room) return; // Forward the parsed payload (sender_platform_user_id / bootstrap_intro) so // callers can ignore completions from other collaborators in shared rooms. const payload = envelope.data && typeof envelope.data === 'object' ? (envelope.data as SuperagentAgentDonePayload) : undefined; handlers.onAgentDone?.(payload); }; const onConnect = () => { emitJoin(room, metadata); handlers.onReconnect?.(); }; const onError: SocketListener = (error) => handlers.onError?.(error); socket.on('update_model', onUpdate); socket.on('agent_done', onAgentDone); socket.on('connect', onConnect); socket.on('error', onError); return () => { removeSocketListener(socket, 'update_model', onUpdate); removeSocketListener(socket, 'agent_done', onAgentDone); removeSocketListener(socket, 'connect', onConnect); removeSocketListener(socket, 'error', onError); removeRoom(room); }; }, }; function addRoom(room: string, metadata?: Record) { const currentCount = roomCounts.get(room) ?? 0; roomCounts.set(room, currentCount + 1); if (currentCount === 0) emitJoin(room, metadata); } function removeRoom(room: string) { const nextCount = (roomCounts.get(room) ?? 1) - 1; if (nextCount > 0) { roomCounts.set(room, nextCount); return; } roomCounts.delete(room); socket.emit('leave', room); } function emitJoin(room: string, metadata?: Record) { if (metadata) socket.emit('join', room, metadata); else socket.emit('join', room); } } function dispatchUpdate(data: unknown, handlers: SuperagentRealtimeHandlers) { if (!data || typeof data !== 'object') return; const payload = data as { _message?: SuperagentMessage; conversation?: SuperagentConversation; id?: unknown; messages?: unknown; }; if (payload._message) handlers.onMessage?.(payload._message); if (payload.conversation) handlers.onConversation?.(payload.conversation); if (typeof payload.id === 'string' && Array.isArray(payload.messages)) { handlers.onConversation?.(payload as SuperagentConversation); } } function parseEnvelope(raw: unknown): RealtimeEnvelope | null { if (!raw || typeof raw !== 'object') return null; const message = raw as { room?: unknown; data?: unknown }; if (typeof message.room !== 'string') return null; return { room: message.room, data: parseData(message.data), }; } function parseData(data: unknown) { if (typeof data !== 'string') return data; try { return JSON.parse(data) as unknown; } catch { return undefined; } } function getConversationRoom(conversationId: string) { const segment = conversationId.trim(); if (!/^[A-Za-z0-9._:-]+$/.test(segment)) { throw new Error('Invalid conversationId'); } return `/agent-conversations/${segment}`; } function removeSocketListener(socket: SuperagentSocketLike, event: string, listener: SocketListener) { if (socket.off) socket.off(event, listener); else socket.removeListener?.(event, listener); }