import WebSocket from 'ws'; import { ClientMessage, safeSend } from '../utils/protocol'; import { isOnline } from '../redis/presenceClient'; import { enqueueOfflineMessage } from '../redis/messageClient'; import { registry } from '../services/connectionRegistry'; import { MessageModel } from '../db/mongo'; import { logger } from '../utils/logger'; const MAX_TEXT_LENGTH = 4096; /** * Handles type: 'message' * * Payload shape: * { * toUid: string, * type: 'text' | 'media', * content: string, // text body or Supabase media URL * mediaType?: string, // MIME type for media * clientMsgId?: string, // Client-generated idempotency ID * } * * Flow: * If recipient online → push directly to their WS * If recipient offline → enqueue in Redis Bucket 2 * Always → persist to MongoDB for history */ export async function handleMessage( senderWs: WebSocket, fromUid: string, msg: ClientMessage ): Promise { const { toUid, type, content, mediaType, clientMsgId } = msg.payload as { toUid?: string; type?: string; content?: string; mediaType?: string; clientMsgId?: string; }; // ── Validate ────────────────────────────────────────────────────────────── if (!toUid || typeof toUid !== 'string') { safeSend(senderWs, 'error', { message: 'toUid is required' }, msg.requestId); return; } if (toUid === fromUid) { safeSend(senderWs, 'error', { message: 'Cannot message yourself' }, msg.requestId); return; } if (!content || typeof content !== 'string') { safeSend(senderWs, 'error', { message: 'content is required' }, msg.requestId); return; } if (type !== 'text' && type !== 'media') { safeSend(senderWs, 'error', { message: 'type must be text or media' }, msg.requestId); return; } if (type === 'text' && content.length > MAX_TEXT_LENGTH) { safeSend(senderWs, 'error', { message: `Text exceeds ${MAX_TEXT_LENGTH} chars` }, msg.requestId); return; } const envelope: Record = { fromUid, toUid, type, content, mediaType: mediaType ?? null, clientMsgId: clientMsgId ?? null, sentAt: Date.now(), }; // ── Persist to MongoDB ──────────────────────────────────────────────────── // Fire-and-forget — do not block message delivery on DB write MessageModel.create({ fromUid, toUid, type, content, mediaType: mediaType ?? undefined, delivered: false, }).catch((err: unknown) => logger.error('MongoDB message persist failed', { error: String(err) }) ); // ── Deliver ─────────────────────────────────────────────────────────────── const targetMeta = registry.getByUid(toUid); if (targetMeta) { // Online — deliver directly safeSend(targetMeta.ws, 'message', envelope); logger.debug('Message delivered live', { fromUid, toUid, type }); } else { // Double-check Redis in case peer is on a different process (future multi-node) const online = await isOnline(toUid).catch(() => false); if (!online) { // Enqueue in Redis Bucket 2 await enqueueOfflineMessage(toUid, envelope).catch((err: unknown) => logger.error('enqueueOfflineMessage failed', { toUid, error: String(err) }) ); logger.debug('Message queued offline', { fromUid, toUid }); } } // Ack the sender safeSend(senderWs, 'message', { ack: true, clientMsgId: clientMsgId ?? null }, msg.requestId); }