chatapi / src /redis /messageClient.ts
wilooper's picture
initial commit
87cb242
Raw History Blame Contribute Delete
3.18 kB
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));
}