File size: 3,838 Bytes
87cb242 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 | 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);
}
|