Download src/handlers/messageHandler.ts from Collabos/chatapi: direct link, hf CLI and curl.
- Browser
- Download file 3.84 kB
-
https://huggingface.co/spaces/Collabos/chatapi/resolve/main/src/handlers/messageHandler.ts
- Command line
-
hf download hf://spaces/Collabos/chatapi/src/handlers/messageHandler.ts
-
curl -L -o messageHandler.ts https://huggingface.co/spaces/Collabos/chatapi/resolve/main/src/handlers/messageHandler.ts
3.84 kB
| 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<void> { | |
| 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<string, unknown> = { | |
| 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); | |
| } | |