import { Queue } from 'bullmq'; import { redis } from './redis'; export const OWNER_QUEUE_NAMES = [ 'moderation', 'backups', 'automod', 'verification', 'reminders', 'giveaways', 'tickets', 'birthdays', 'stats', 'feeds', 'schedules', 'guild-backups', 'suggestions', 'presence', 'retention' ] as const; export type OwnerQueueName = (typeof OWNER_QUEUE_NAMES)[number]; const globalForOwnerQueues = globalThis as unknown as { __nexumiOwnerQueues?: Map; }; function getQueue(name: string): Queue { const cache = globalForOwnerQueues.__nexumiOwnerQueues ?? new Map(); globalForOwnerQueues.__nexumiOwnerQueues = cache; let queue = cache.get(name); if (!queue) { queue = new Queue(name, { connection: redis }); cache.set(name, queue); } return queue; } export interface QueueStats { name: string; waiting: number; active: number; completed: number; failed: number; delayed: number; } export async function getAllQueueStats(): Promise { return Promise.all( OWNER_QUEUE_NAMES.map(async (name) => { const queue = getQueue(name); const counts = await queue.getJobCounts('wait', 'active', 'completed', 'failed', 'delayed'); return { name, waiting: counts.wait, active: counts.active, completed: counts.completed, failed: counts.failed, delayed: counts.delayed }; }) ); } export async function retryFailedJobs(queueName: string, limit = 25): Promise { if (!(OWNER_QUEUE_NAMES as readonly string[]).includes(queueName)) { throw new Error('Unknown queue'); } const queue = getQueue(queueName); const failed = await queue.getFailed(0, limit - 1); let retried = 0; for (const job of failed) { await job.retry(); retried += 1; } return retried; }