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 ): Promise { 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>> { 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> = []; for (const raw of rawList) { try { parsed.push(JSON.parse(raw) as Record); } catch { // Corrupted entry — skip } } return parsed; } /** How many messages are queued for a user. Useful for dashboards. */ export async function getQueueLength(uid: string): Promise { return messageClient.llen(queueKey(uid)); }