File size: 3,178 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
import Redis from 'ioredis';
import { CONFIG } from '../config';
import { logger } from '../utils/logger';

/**
 * Redis Bucket 2 β€” Offline Message Queue
 *
 * Key schema:
 *   offline_msgs:{uid}  β†’ Redis List of JSON-encoded message strings
 *                          Each list key has OFFLINE_MSG_TTL_S TTL.
 *
 * Entirely separate Redis instance from Bucket 1.
 */
const messageClient = new Redis(CONFIG.REDIS_MESSAGE_URL, {
  maxRetriesPerRequest: 2,
  connectTimeout: 5000,
  lazyConnect: false,
  enableReadyCheck: true,
  keepAlive: 10000,
});

messageClient.on('connect', () => logger.info('Message Redis connected'));
messageClient.on('error', (err: Error) =>
  logger.error('Message Redis error', { error: err.message })
);
messageClient.on('reconnecting', () => logger.warn('Message Redis reconnecting'));

export { messageClient };

// ── Key builder ──────────────────────────────────────────────────────────────

const queueKey = (uid: string) => `offline_msgs:${uid}`;

// Max messages to queue per user β€” prevents memory abuse at 10k scale
const MAX_QUEUE_LENGTH = 200;

// ── Operations ───────────────────────────────────────────────────────────────

/** Push a serialized message onto the offline queue for a user. */
export async function enqueueOfflineMessage(
  toUid: string,
  message: Record<string, unknown>
): Promise<void> {
  const key = queueKey(toUid);
  const serialized = JSON.stringify(message);

  // RPUSH then cap the list, then refresh TTL β€” 3 commands in a pipeline
  const pipeline = messageClient.pipeline();
  pipeline.rpush(key, serialized);
  // Trim to MAX_QUEUE_LENGTH (drop oldest if overflow)
  pipeline.ltrim(key, -MAX_QUEUE_LENGTH, -1);
  pipeline.expire(key, CONFIG.OFFLINE_MSG_TTL_S);
  await pipeline.exec();
}

/**
 * Atomically pop all queued messages for a uid.
 * Uses LRANGE + DEL so no message is double-delivered.
 * Returns parsed message objects.
 */
export async function flushOfflineQueue(
  toUid: string
): Promise<Array<Record<string, unknown>>> {
  const key = queueKey(toUid);

  // LRANGE 0 -1 gets everything; DEL removes the key atomically-ish
  // For true atomicity at high scale, this could be a Lua script β€”
  // but for HF free tier the risk of race is negligible.
  const pipeline = messageClient.pipeline();
  pipeline.lrange(key, 0, -1);
  pipeline.del(key);
  const results = await pipeline.exec();

  if (!results || !results[0] || results[0][0]) return [];

  const rawList = results[0][1] as string[];
  const parsed: Array<Record<string, unknown>> = [];

  for (const raw of rawList) {
    try {
      parsed.push(JSON.parse(raw) as Record<string, unknown>);
    } catch {
      // Corrupted entry β€” skip
    }
  }

  return parsed;
}

/** How many messages are queued for a user. Useful for dashboards. */
export async function getQueueLength(uid: string): Promise<number> {
  return messageClient.llen(queueKey(uid));
}