diff --git a/src/modules/announcements/announcement.module.ts b/src/modules/announcements/announcement.module.ts index 255a778..5d2396f 100755 --- a/src/modules/announcements/announcement.module.ts +++ b/src/modules/announcements/announcement.module.ts @@ -12,13 +12,14 @@ import { DanakServicesModule } from "../danak-services/danak-services.module"; import { UsersModule } from "../users/users.module"; import { UserAnnouncement } from "./entities/user-announcement.entity"; import { UserAnnouncementRepository } from "./repositories/user-announcement.repository"; - +import { NotificationModule } from "../notifications/notifications.module"; @Module({ imports: [ TypeOrmModule.forFeature([Announcement, UserAnnouncement]), BullModule.registerQueue({ name: ANNOUNCEMENT.ANNOUNCEMENT_QUEUE_NAME }), DanakServicesModule, UsersModule, + NotificationModule, ], providers: [AnnouncementService, AnnouncementRepository, UserAnnouncementRepository, AnnouncementProcessor], controllers: [AnnouncementController], diff --git a/src/modules/announcements/queue/announcment.processor.ts b/src/modules/announcements/queue/announcment.processor.ts index eb678e4..8f02bef 100755 --- a/src/modules/announcements/queue/announcment.processor.ts +++ b/src/modules/announcements/queue/announcment.processor.ts @@ -4,17 +4,25 @@ import { Job } from "bullmq"; import { DataSource } from "typeorm"; import { WorkerProcessor } from "../../../common/queues/worker.processor"; +import { IAnnouncementNotificationData } from "../../notifications/interfaces/ISendNotificationData"; +import { NotificationQueue } from "../../notifications/queue/notification.queue"; import { SubscriptionPlan } from "../../subscriptions/entities/subscription.entity"; import { UserSubscription } from "../../subscriptions/entities/user-subscription.entity"; import { User } from "../../users/entities/user.entity"; +import { UsersService } from "../../users/providers/users.service"; import { ANNOUNCEMENT } from "../constants"; +import { Announcement } from "../entities/announcement.entity"; import { UserAnnouncement } from "../entities/user-announcement.entity"; import { ISendAnnouncement } from "../interfaces/ISendAnnouncement"; @Processor(ANNOUNCEMENT.ANNOUNCEMENT_QUEUE_NAME) export class AnnouncementProcessor extends WorkerProcessor { protected readonly logger = new Logger(AnnouncementProcessor.name); - constructor(private readonly dataSource: DataSource) { + constructor( + private readonly dataSource: DataSource, + private readonly notificationQueue: NotificationQueue, + private readonly usersService: UsersService, + ) { super(); } @@ -23,7 +31,7 @@ export class AnnouncementProcessor extends WorkerProcessor { switch (job.name) { case ANNOUNCEMENT.ANNOUNCEMENT_SEND_JOB_NAME: - return this.sendAnnouncement(job, token); + return await this.sendAnnouncement(job, token); case ANNOUNCEMENT.ANNOUNCEMENT_PUBLISH_JOB_NAME: return this.publishAnnouncement(job, token); default: @@ -36,11 +44,12 @@ export class AnnouncementProcessor extends WorkerProcessor { this.logger.log(`Sending announcement: ${job.data.announcementId} to users. ${token}`); const queryRunner = this.dataSource.createQueryRunner(); - await queryRunner.connect(); - await queryRunner.startTransaction(); + const data = job.data; try { + await queryRunner.connect(); + await queryRunner.startTransaction(); // Create a query builder to fetch unique users const userQueryBuilder = queryRunner.manager.createQueryBuilder(User, "user").select("user.id").distinct(true); @@ -89,8 +98,26 @@ export class AnnouncementProcessor extends WorkerProcessor { await queryRunner.manager.insert(UserAnnouncement, userAnnouncements); + // Fetch announcement details for notification + const announcement = await queryRunner.manager.findOne(Announcement, { where: { id: job.data.announcementId } }); + if (!announcement) throw new Error("Announcement not found"); + + // Send notification to each user + for (const userId of users) { + const user = await this.usersService.findOneByIdWithQueryRunner(userId, queryRunner); + const notifData: IAnnouncementNotificationData = { + userPhone: user.phone, + userEmail: user.email, + title: announcement.title, + description: announcement.content, + date: announcement.publishAt ?? new Date(), + }; + await this.notificationQueue.addAnnouncementNotification(userId, notifData); + } + this.logger.log(`Successfully sent announcement ${announcementId} to ${users.length} users`); await queryRunner.commitTransaction(); + return true; } catch (error) { this.logger.error(`Failed to send announcement:`, error); await queryRunner.rollbackTransaction(); diff --git a/src/modules/notifications/constants/index.ts b/src/modules/notifications/constants/index.ts index 46c177c..199084c 100644 --- a/src/modules/notifications/constants/index.ts +++ b/src/modules/notifications/constants/index.ts @@ -7,4 +7,9 @@ export const NOTIFICATION = Object.freeze({ SEND_NOTIFICATION_JOB_ATTEMPTS: 3, // retry 3 times SEND_NOTIFICATION_JOB_BACKOFF: 5 * 1000, // retry after 5 seconds SEND_NOTIFICATION_JOB_TIMEOUT: 10000, // timeout after 10 seconds + + SEND_NOTIFICATION_JOB_CONCURRENCY: 2, + SEND_NOTIFICATION_JOB_LOCK_DURATION: 30000, + SEND_NOTIFICATION_JOB_STALLED_INTERVAL: 30000, + SEND_NOTIFICATION_JOB_MAX_STALLED_COUNT: 2, }); diff --git a/src/modules/notifications/queue/notification.processor.ts b/src/modules/notifications/queue/notification.processor.ts index 0992ab4..184ac9f 100644 --- a/src/modules/notifications/queue/notification.processor.ts +++ b/src/modules/notifications/queue/notification.processor.ts @@ -42,10 +42,10 @@ type NotificationJobData = { }; @Processor(NOTIFICATION.QUEUE_NAME, { - concurrency: 5, - lockDuration: 30000, // 30 seconds lock duration - stalledInterval: 30000, // Check for stalled jobs every 30 seconds - maxStalledCount: 2, // Allow 2 stalls before failing + concurrency: NOTIFICATION.SEND_NOTIFICATION_JOB_CONCURRENCY, + lockDuration: NOTIFICATION.SEND_NOTIFICATION_JOB_LOCK_DURATION, + stalledInterval: NOTIFICATION.SEND_NOTIFICATION_JOB_STALLED_INTERVAL, + maxStalledCount: NOTIFICATION.SEND_NOTIFICATION_JOB_MAX_STALLED_COUNT, }) export class NotificationProcessor extends WorkerProcessor { protected readonly logger = new Logger(NotificationProcessor.name); @@ -60,136 +60,105 @@ export class NotificationProcessor extends WorkerProcessor { async process(job: Job, token?: string) { this.logger.log(`Processing notification job: ${job.id} ${token ? `with token: ${token}` : ""}`); - const queryRunner = this.dataSource.createQueryRunner(); + let heartbeat: NodeJS.Timeout | undefined; + + // try { await queryRunner.connect(); await queryRunner.startTransaction(); + heartbeat = setInterval(() => job.updateProgress(50), 10000); - const heartbeat = setInterval(() => { - job.updateProgress(50); - }, 10000); // Update progress every 10 seconds - + // const { type, recipientId, data } = job.data; - - try { - switch (type) { - // Admin Notifications - case NotifType.NEW_BLOG_COMMENT: - await this.notificationsService.createNewBlogCommentNotification( - recipientId, - data as IBlogCommentNotificationData, - queryRunner, - ); - break; - case NotifType.NEW_SERVICE_REVIEW: - await this.notificationsService.createNewServiceReviewNotification( - recipientId, - data as IServiceReviewNotificationData, - queryRunner, - ); - break; - case NotifType.NEW_CUSTOMER: - await this.notificationsService.createNewCustomerNotification(recipientId, data as INewCustomerNotificationData, queryRunner); - break; - case NotifType.NEW_SUBSCRIPTION: - await this.notificationsService.createNewSubscriptionNotification( - recipientId, - data as INewSubscriptionNotificationData, - queryRunner, - ); - break; - case NotifType.NEW_CRITICISM: - await this.notificationsService.createNewCriticismNotification(recipientId, data as INewCriticismNotificationData, queryRunner); - break; - case NotifType.NEW_TICKET: - await this.notificationsService.createNewTicketGlobalNotification(recipientId, data as INewTicketNotificationData, queryRunner); - break; - - // User Notifications - case NotifType.USER_LOGIN: - await this.notificationsService.createLoginNotification(recipientId, data); - break; - case NotifType.ANNOUNCEMENT: - await this.notificationsService.createAnnouncementNotification(recipientId, data as IAnnouncementNotificationData, queryRunner); - break; - - // Wallet Notifications - case NotifType.WALLET_CHARGE: - await this.notificationsService.createWalletChargeNotification(recipientId, data as IWalletNotificationData, queryRunner); - break; - case NotifType.WALLET_DEDUCTION: - await this.notificationsService.createWalletDeductionNotification(recipientId, data as IWalletNotificationData, queryRunner); - break; - - // Ticket Notifications - case NotifType.ANSWER_TICKET: - await this.notificationsService.createAnswerTicketNotification(recipientId, data as ITicketNotificationData, queryRunner); - break; - case NotifType.CREATE_TICKET: - await this.notificationsService.createTicketNotification(recipientId, data as ITicketNotificationData, queryRunner); - break; - case NotifType.ASSIGN_TICKET: - await this.notificationsService.createAssignTicketNotificationForAdmin( - recipientId, - data as ITicketNotificationData, - queryRunner, - ); - break; - - // Invoice Notifications - case NotifType.CREATE_INVOICE: - await this.notificationsService.createInvoiceCreationNotification(recipientId, data as IInvoiceNotificationData, queryRunner); - break; - case NotifType.BILL_INVOICE_REMINDER: - await this.notificationsService.createBillInvoiceReminderNotification( - recipientId, - data as IInvoiceNotificationData, - queryRunner, - ); - break; - case NotifType.BILL_INVOICE: - await this.notificationsService.createBillInvoiceNotification(recipientId, data as IInvoiceNotificationData, queryRunner); - break; - case NotifType.APPROVE_INVOICE: - await this.notificationsService.createApprovedInvoiceNotification(recipientId, data as IInvoiceNotificationData, queryRunner); - break; - case NotifType.INVOICE_OVERDUE: - await this.notificationsService.createInvoiceOverdueNotification(recipientId, data as IInvoiceNotificationData, queryRunner); - break; - case NotifType.RECURRING_INVOICE: - await this.notificationsService.createRecurringInvoiceNotification(recipientId, data as IInvoiceNotificationData, queryRunner); - break; - - // Service Notifications - case NotifType.BLOCK_SERVICE: - await this.notificationsService.createBlockServiceNotification(recipientId, data as ISubscriptionNotificationData, queryRunner); - break; - - // Payment Notifications - case NotifType.PAYMENT_REMINDER: - await this.notificationsService.createPaymentReminderNotification(recipientId, data as IPaymentNotificationData, queryRunner); - break; - case NotifType.PAYMENT_CANCELLATION: - await this.notificationsService.createPaymentCancellationNotification( - recipientId, - data as IPaymentNotificationData, - queryRunner, - ); - break; - - default: - this.logger.warn(`Unknown notification type: ${type}`); - } - - await queryRunner.commitTransaction(); - clearInterval(heartbeat); - return true; - } catch (error) { - clearInterval(heartbeat); - throw error; + switch (type) { + // Admin Notifications + case NotifType.NEW_BLOG_COMMENT: + await this.notificationsService.createNewBlogCommentNotification(recipientId, data as IBlogCommentNotificationData, queryRunner); + break; + case NotifType.NEW_SERVICE_REVIEW: + await this.notificationsService.createNewServiceReviewNotification( + recipientId, + data as IServiceReviewNotificationData, + queryRunner, + ); + break; + case NotifType.NEW_CUSTOMER: + await this.notificationsService.createNewCustomerNotification(recipientId, data as INewCustomerNotificationData, queryRunner); + break; + case NotifType.NEW_SUBSCRIPTION: + await this.notificationsService.createNewSubscriptionNotification( + recipientId, + data as INewSubscriptionNotificationData, + queryRunner, + ); + break; + case NotifType.NEW_CRITICISM: + await this.notificationsService.createNewCriticismNotification(recipientId, data as INewCriticismNotificationData, queryRunner); + break; + case NotifType.NEW_TICKET: + await this.notificationsService.createNewTicketGlobalNotification(recipientId, data as INewTicketNotificationData, queryRunner); + break; + // User Notifications + case NotifType.USER_LOGIN: + await this.notificationsService.createLoginNotification(recipientId, data); + break; + case NotifType.ANNOUNCEMENT: + await this.notificationsService.createAnnouncementNotification(recipientId, data as IAnnouncementNotificationData, queryRunner); + break; + // Wallet Notifications + case NotifType.WALLET_CHARGE: + await this.notificationsService.createWalletChargeNotification(recipientId, data as IWalletNotificationData, queryRunner); + break; + case NotifType.WALLET_DEDUCTION: + await this.notificationsService.createWalletDeductionNotification(recipientId, data as IWalletNotificationData, queryRunner); + break; + // Ticket Notifications + case NotifType.ANSWER_TICKET: + await this.notificationsService.createAnswerTicketNotification(recipientId, data as ITicketNotificationData, queryRunner); + break; + case NotifType.CREATE_TICKET: + await this.notificationsService.createTicketNotification(recipientId, data as ITicketNotificationData, queryRunner); + break; + case NotifType.ASSIGN_TICKET: + await this.notificationsService.createAssignTicketNotificationForAdmin(recipientId, data as ITicketNotificationData, queryRunner); + break; + // Invoice Notifications + case NotifType.CREATE_INVOICE: + await this.notificationsService.createInvoiceCreationNotification(recipientId, data as IInvoiceNotificationData, queryRunner); + break; + case NotifType.BILL_INVOICE_REMINDER: + await this.notificationsService.createBillInvoiceReminderNotification(recipientId, data as IInvoiceNotificationData, queryRunner); + break; + case NotifType.BILL_INVOICE: + await this.notificationsService.createBillInvoiceNotification(recipientId, data as IInvoiceNotificationData, queryRunner); + break; + case NotifType.APPROVE_INVOICE: + await this.notificationsService.createApprovedInvoiceNotification(recipientId, data as IInvoiceNotificationData, queryRunner); + break; + case NotifType.INVOICE_OVERDUE: + await this.notificationsService.createInvoiceOverdueNotification(recipientId, data as IInvoiceNotificationData, queryRunner); + break; + case NotifType.RECURRING_INVOICE: + await this.notificationsService.createRecurringInvoiceNotification(recipientId, data as IInvoiceNotificationData, queryRunner); + break; + // Service Notifications + case NotifType.BLOCK_SERVICE: + await this.notificationsService.createBlockServiceNotification(recipientId, data as ISubscriptionNotificationData, queryRunner); + break; + // Payment Notifications + case NotifType.PAYMENT_REMINDER: + await this.notificationsService.createPaymentReminderNotification(recipientId, data as IPaymentNotificationData, queryRunner); + break; + case NotifType.PAYMENT_CANCELLATION: + await this.notificationsService.createPaymentCancellationNotification(recipientId, data as IPaymentNotificationData, queryRunner); + break; + default: + this.logger.warn(`Unknown notification type: ${type}`); } + await queryRunner.commitTransaction(); + return true; } catch (error) { this.logger.error( `Failed to process notification: ${error instanceof Error ? error.message : "Unknown error"}`, @@ -198,6 +167,7 @@ export class NotificationProcessor extends WorkerProcessor { await queryRunner.rollbackTransaction(); throw error; } finally { + if (heartbeat) clearInterval(heartbeat); await queryRunner.release(); } } diff --git a/src/modules/notifications/queue/notification.queue.ts b/src/modules/notifications/queue/notification.queue.ts index c627f59..8144ff3 100644 --- a/src/modules/notifications/queue/notification.queue.ts +++ b/src/modules/notifications/queue/notification.queue.ts @@ -87,7 +87,6 @@ export class NotificationQueue { async addWalletChargeNotification(recipientId: string, data: IWalletNotificationData) { await this.addNotificationJob(NotifType.WALLET_CHARGE, recipientId, data); } - //TODO:USER THIS async addWalletDeductionNotification(recipientId: string, data: IWalletNotificationData) { await this.addNotificationJob(NotifType.WALLET_DEDUCTION, recipientId, data);