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);
}