Download src/redis/messageClient.ts from Collabos/chatapi: direct link, hf CLI and curl.
- Browser
- Download file 3.18 kB
-
https://huggingface.co/spaces/Collabos/chatapi/resolve/main/src/redis/messageClient.ts
- Command line
-
hf download hf://spaces/Collabos/chatapi/src/redis/messageClient.ts
-
curl -L -o messageClient.ts https://huggingface.co/spaces/Collabos/chatapi/resolve/main/src/redis/messageClient.ts
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)); | |
| } | |