import { Injectable, Logger } from '@nestjs/common'; export interface QueuedMessage { id: string; levelId: number; levelType: string; recipient: string; message: string; type?: 'EMAIL' | 'SMS' | 'WA' | 'TELEPHONE'; priority: 'high' | 'medium' | 'low'; scheduledFor?: Date; metadata?: any; createdAt: Date; retryCount: number; maxRetries: number; } @Injectable() export class IntegrationQueueService { private readonly logger = new Logger(IntegrationQueueService.name); private queue: QueuedMessage[] = []; private processing = false; constructor() { // Start processing queue every 5 seconds setInterval(() => this.processQueue(), 5000); } async queueMessage( levelId: number, levelType: string, recipient: string, message: string, type?: 'EMAIL' | 'SMS' | 'WA' | 'TELEPHONE', priority: 'high' | 'medium' | 'low' = 'medium', scheduledFor?: Date, metadata?: any, ): Promise { const messageId = this.generateMessageId(); const queuedMessage: QueuedMessage = { id: messageId, levelId, levelType, recipient, message, type, priority, scheduledFor, metadata, createdAt: new Date(), retryCount: 0, maxRetries: 3, }; this.queue.push(queuedMessage); this.sortQueue(); this.logger.log( `Message queued: ${messageId} for ${recipient} with priority ${priority}`, ); return messageId; } async queueBulkMessages( levelId: number, levelType: string, recipients: string[], message: string, type?: 'EMAIL' | 'SMS' | 'WA' | 'TELEPHONE', priority: 'high' | 'medium' | 'low' = 'low', metadata?: any, ): Promise { const messageIds: string[] = []; for (const recipient of recipients) { const messageId = await this.queueMessage( levelId, levelType, recipient, message, type, priority, undefined, metadata, ); messageIds.push(messageId); } return messageIds; } getQueueStats() { const total = this.queue.length; const pending = this.queue.filter( (m) => !m.scheduledFor || m.scheduledFor <= new Date(), ).length; const scheduled = this.queue.filter( (m) => m.scheduledFor && m.scheduledFor > new Date(), ).length; const byPriority = { high: this.queue.filter((m) => m.priority === 'high').length, medium: this.queue.filter((m) => m.priority === 'medium').length, low: this.queue.filter((m) => m.priority === 'low').length, }; return { total, pending, scheduled, byPriority, processing: this.processing, }; } private async processQueue(): Promise { if (this.processing || this.queue.length === 0) { return; } this.processing = true; try { const now = new Date(); // Get messages ready to send (not scheduled or scheduled time has passed) const readyMessages = this.queue.filter( (message) => !message.scheduledFor || message.scheduledFor <= now, ); if (readyMessages.length === 0) { return; } // Sort by priority (high -> medium -> low) const prioritySorted = readyMessages.sort((a, b) => { const priorityOrder = { high: 3, medium: 2, low: 1 }; return priorityOrder[b.priority] - priorityOrder[a.priority]; }); // Process up to 5 messages at once const messagesToProcess = prioritySorted.slice(0, 5); for (const message of messagesToProcess) { try { // Remove message from queue this.queue = this.queue.filter((m) => m.id !== message.id); this.logger.log( `Processing queued message ${message.id} for ${message.recipient}`, ); // Note: In a real implementation, you would call the communication service here // For now, we just simulate successful processing } catch (error) { this.logger.error( `Failed to process message ${message.id}:`, error.message, ); // Retry logic if (message.retryCount < message.maxRetries) { message.retryCount++; message.scheduledFor = new Date( Date.now() + message.retryCount * 60000, ); // Retry after 1, 2, 3 minutes this.queue.push(message); this.logger.warn( `Message ${message.id} will be retried (attempt ${message.retryCount}/${message.maxRetries})`, ); } else { this.logger.error( `Message ${message.id} failed permanently after ${message.maxRetries} attempts`, ); } } } } catch (error) { this.logger.error('Error processing queue:', error.message); } finally { this.processing = false; } } private sortQueue(): void { this.queue.sort((a, b) => { // Sort by priority first const priorityOrder = { high: 3, medium: 2, low: 1 }; if (priorityOrder[a.priority] !== priorityOrder[b.priority]) { return priorityOrder[b.priority] - priorityOrder[a.priority]; } // Then by creation time (oldest first) return a.createdAt.getTime() - b.createdAt.getTime(); }); } private generateMessageId(): string { return `msg_${Date.now()}_${Math.random().toString(36).substr(2, 9)}`; } // Helper methods for integration isQueueEmpty(): boolean { return this.queue.length === 0; } getQueueLength(): number { return this.queue.length; } clearQueue(): void { this.queue = []; this.logger.log('Queue cleared'); } // Cancel a specific message cancelMessage(messageId: string): boolean { const initialLength = this.queue.length; this.queue = this.queue.filter((m) => m.id !== messageId); const canceled = this.queue.length < initialLength; if (canceled) { this.logger.log(`Message ${messageId} canceled`); } return canceled; } }