chore: add send announce sms

This commit is contained in:
Mahyargdz
2025-05-14 10:03:51 +03:30
parent d2af460de3
commit 9bd662f872
5 changed files with 134 additions and 132 deletions
@@ -12,13 +12,14 @@ import { DanakServicesModule } from "../danak-services/danak-services.module";
import { UsersModule } from "../users/users.module"; import { UsersModule } from "../users/users.module";
import { UserAnnouncement } from "./entities/user-announcement.entity"; import { UserAnnouncement } from "./entities/user-announcement.entity";
import { UserAnnouncementRepository } from "./repositories/user-announcement.repository"; import { UserAnnouncementRepository } from "./repositories/user-announcement.repository";
import { NotificationModule } from "../notifications/notifications.module";
@Module({ @Module({
imports: [ imports: [
TypeOrmModule.forFeature([Announcement, UserAnnouncement]), TypeOrmModule.forFeature([Announcement, UserAnnouncement]),
BullModule.registerQueue({ name: ANNOUNCEMENT.ANNOUNCEMENT_QUEUE_NAME }), BullModule.registerQueue({ name: ANNOUNCEMENT.ANNOUNCEMENT_QUEUE_NAME }),
DanakServicesModule, DanakServicesModule,
UsersModule, UsersModule,
NotificationModule,
], ],
providers: [AnnouncementService, AnnouncementRepository, UserAnnouncementRepository, AnnouncementProcessor], providers: [AnnouncementService, AnnouncementRepository, UserAnnouncementRepository, AnnouncementProcessor],
controllers: [AnnouncementController], controllers: [AnnouncementController],
@@ -4,17 +4,25 @@ import { Job } from "bullmq";
import { DataSource } from "typeorm"; import { DataSource } from "typeorm";
import { WorkerProcessor } from "../../../common/queues/worker.processor"; 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 { SubscriptionPlan } from "../../subscriptions/entities/subscription.entity";
import { UserSubscription } from "../../subscriptions/entities/user-subscription.entity"; import { UserSubscription } from "../../subscriptions/entities/user-subscription.entity";
import { User } from "../../users/entities/user.entity"; import { User } from "../../users/entities/user.entity";
import { UsersService } from "../../users/providers/users.service";
import { ANNOUNCEMENT } from "../constants"; import { ANNOUNCEMENT } from "../constants";
import { Announcement } from "../entities/announcement.entity";
import { UserAnnouncement } from "../entities/user-announcement.entity"; import { UserAnnouncement } from "../entities/user-announcement.entity";
import { ISendAnnouncement } from "../interfaces/ISendAnnouncement"; import { ISendAnnouncement } from "../interfaces/ISendAnnouncement";
@Processor(ANNOUNCEMENT.ANNOUNCEMENT_QUEUE_NAME) @Processor(ANNOUNCEMENT.ANNOUNCEMENT_QUEUE_NAME)
export class AnnouncementProcessor extends WorkerProcessor { export class AnnouncementProcessor extends WorkerProcessor {
protected readonly logger = new Logger(AnnouncementProcessor.name); 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(); super();
} }
@@ -23,7 +31,7 @@ export class AnnouncementProcessor extends WorkerProcessor {
switch (job.name) { switch (job.name) {
case ANNOUNCEMENT.ANNOUNCEMENT_SEND_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: case ANNOUNCEMENT.ANNOUNCEMENT_PUBLISH_JOB_NAME:
return this.publishAnnouncement(job, token); return this.publishAnnouncement(job, token);
default: default:
@@ -36,11 +44,12 @@ export class AnnouncementProcessor extends WorkerProcessor {
this.logger.log(`Sending announcement: ${job.data.announcementId} to users. ${token}`); this.logger.log(`Sending announcement: ${job.data.announcementId} to users. ${token}`);
const queryRunner = this.dataSource.createQueryRunner(); const queryRunner = this.dataSource.createQueryRunner();
await queryRunner.connect();
await queryRunner.startTransaction();
const data = job.data; const data = job.data;
try { try {
await queryRunner.connect();
await queryRunner.startTransaction();
// Create a query builder to fetch unique users // Create a query builder to fetch unique users
const userQueryBuilder = queryRunner.manager.createQueryBuilder(User, "user").select("user.id").distinct(true); 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); 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`); this.logger.log(`Successfully sent announcement ${announcementId} to ${users.length} users`);
await queryRunner.commitTransaction(); await queryRunner.commitTransaction();
return true;
} catch (error) { } catch (error) {
this.logger.error(`Failed to send announcement:`, error); this.logger.error(`Failed to send announcement:`, error);
await queryRunner.rollbackTransaction(); await queryRunner.rollbackTransaction();
@@ -7,4 +7,9 @@ export const NOTIFICATION = Object.freeze({
SEND_NOTIFICATION_JOB_ATTEMPTS: 3, // retry 3 times SEND_NOTIFICATION_JOB_ATTEMPTS: 3, // retry 3 times
SEND_NOTIFICATION_JOB_BACKOFF: 5 * 1000, // retry after 5 seconds SEND_NOTIFICATION_JOB_BACKOFF: 5 * 1000, // retry after 5 seconds
SEND_NOTIFICATION_JOB_TIMEOUT: 10000, // timeout after 10 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,
}); });
@@ -42,10 +42,10 @@ type NotificationJobData = {
}; };
@Processor(NOTIFICATION.QUEUE_NAME, { @Processor(NOTIFICATION.QUEUE_NAME, {
concurrency: 5, concurrency: NOTIFICATION.SEND_NOTIFICATION_JOB_CONCURRENCY,
lockDuration: 30000, // 30 seconds lock duration lockDuration: NOTIFICATION.SEND_NOTIFICATION_JOB_LOCK_DURATION,
stalledInterval: 30000, // Check for stalled jobs every 30 seconds stalledInterval: NOTIFICATION.SEND_NOTIFICATION_JOB_STALLED_INTERVAL,
maxStalledCount: 2, // Allow 2 stalls before failing maxStalledCount: NOTIFICATION.SEND_NOTIFICATION_JOB_MAX_STALLED_COUNT,
}) })
export class NotificationProcessor extends WorkerProcessor { export class NotificationProcessor extends WorkerProcessor {
protected readonly logger = new Logger(NotificationProcessor.name); protected readonly logger = new Logger(NotificationProcessor.name);
@@ -60,28 +60,22 @@ export class NotificationProcessor extends WorkerProcessor {
async process(job: Job<NotificationJobData>, token?: string) { async process(job: Job<NotificationJobData>, token?: string) {
this.logger.log(`Processing notification job: ${job.id} ${token ? `with token: ${token}` : ""}`); this.logger.log(`Processing notification job: ${job.id} ${token ? `with token: ${token}` : ""}`);
const queryRunner = this.dataSource.createQueryRunner(); const queryRunner = this.dataSource.createQueryRunner();
let heartbeat: NodeJS.Timeout | undefined;
//
try { try {
await queryRunner.connect(); await queryRunner.connect();
await queryRunner.startTransaction(); 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; const { type, recipientId, data } = job.data;
try {
switch (type) { switch (type) {
// Admin Notifications // Admin Notifications
case NotifType.NEW_BLOG_COMMENT: case NotifType.NEW_BLOG_COMMENT:
await this.notificationsService.createNewBlogCommentNotification( await this.notificationsService.createNewBlogCommentNotification(recipientId, data as IBlogCommentNotificationData, queryRunner);
recipientId,
data as IBlogCommentNotificationData,
queryRunner,
);
break; break;
case NotifType.NEW_SERVICE_REVIEW: case NotifType.NEW_SERVICE_REVIEW:
await this.notificationsService.createNewServiceReviewNotification( await this.notificationsService.createNewServiceReviewNotification(
@@ -106,7 +100,6 @@ export class NotificationProcessor extends WorkerProcessor {
case NotifType.NEW_TICKET: case NotifType.NEW_TICKET:
await this.notificationsService.createNewTicketGlobalNotification(recipientId, data as INewTicketNotificationData, queryRunner); await this.notificationsService.createNewTicketGlobalNotification(recipientId, data as INewTicketNotificationData, queryRunner);
break; break;
// User Notifications // User Notifications
case NotifType.USER_LOGIN: case NotifType.USER_LOGIN:
await this.notificationsService.createLoginNotification(recipientId, data); await this.notificationsService.createLoginNotification(recipientId, data);
@@ -114,7 +107,6 @@ export class NotificationProcessor extends WorkerProcessor {
case NotifType.ANNOUNCEMENT: case NotifType.ANNOUNCEMENT:
await this.notificationsService.createAnnouncementNotification(recipientId, data as IAnnouncementNotificationData, queryRunner); await this.notificationsService.createAnnouncementNotification(recipientId, data as IAnnouncementNotificationData, queryRunner);
break; break;
// Wallet Notifications // Wallet Notifications
case NotifType.WALLET_CHARGE: case NotifType.WALLET_CHARGE:
await this.notificationsService.createWalletChargeNotification(recipientId, data as IWalletNotificationData, queryRunner); await this.notificationsService.createWalletChargeNotification(recipientId, data as IWalletNotificationData, queryRunner);
@@ -122,7 +114,6 @@ export class NotificationProcessor extends WorkerProcessor {
case NotifType.WALLET_DEDUCTION: case NotifType.WALLET_DEDUCTION:
await this.notificationsService.createWalletDeductionNotification(recipientId, data as IWalletNotificationData, queryRunner); await this.notificationsService.createWalletDeductionNotification(recipientId, data as IWalletNotificationData, queryRunner);
break; break;
// Ticket Notifications // Ticket Notifications
case NotifType.ANSWER_TICKET: case NotifType.ANSWER_TICKET:
await this.notificationsService.createAnswerTicketNotification(recipientId, data as ITicketNotificationData, queryRunner); await this.notificationsService.createAnswerTicketNotification(recipientId, data as ITicketNotificationData, queryRunner);
@@ -131,23 +122,14 @@ export class NotificationProcessor extends WorkerProcessor {
await this.notificationsService.createTicketNotification(recipientId, data as ITicketNotificationData, queryRunner); await this.notificationsService.createTicketNotification(recipientId, data as ITicketNotificationData, queryRunner);
break; break;
case NotifType.ASSIGN_TICKET: case NotifType.ASSIGN_TICKET:
await this.notificationsService.createAssignTicketNotificationForAdmin( await this.notificationsService.createAssignTicketNotificationForAdmin(recipientId, data as ITicketNotificationData, queryRunner);
recipientId,
data as ITicketNotificationData,
queryRunner,
);
break; break;
// Invoice Notifications // Invoice Notifications
case NotifType.CREATE_INVOICE: case NotifType.CREATE_INVOICE:
await this.notificationsService.createInvoiceCreationNotification(recipientId, data as IInvoiceNotificationData, queryRunner); await this.notificationsService.createInvoiceCreationNotification(recipientId, data as IInvoiceNotificationData, queryRunner);
break; break;
case NotifType.BILL_INVOICE_REMINDER: case NotifType.BILL_INVOICE_REMINDER:
await this.notificationsService.createBillInvoiceReminderNotification( await this.notificationsService.createBillInvoiceReminderNotification(recipientId, data as IInvoiceNotificationData, queryRunner);
recipientId,
data as IInvoiceNotificationData,
queryRunner,
);
break; break;
case NotifType.BILL_INVOICE: case NotifType.BILL_INVOICE:
await this.notificationsService.createBillInvoiceNotification(recipientId, data as IInvoiceNotificationData, queryRunner); await this.notificationsService.createBillInvoiceNotification(recipientId, data as IInvoiceNotificationData, queryRunner);
@@ -161,35 +143,22 @@ export class NotificationProcessor extends WorkerProcessor {
case NotifType.RECURRING_INVOICE: case NotifType.RECURRING_INVOICE:
await this.notificationsService.createRecurringInvoiceNotification(recipientId, data as IInvoiceNotificationData, queryRunner); await this.notificationsService.createRecurringInvoiceNotification(recipientId, data as IInvoiceNotificationData, queryRunner);
break; break;
// Service Notifications // Service Notifications
case NotifType.BLOCK_SERVICE: case NotifType.BLOCK_SERVICE:
await this.notificationsService.createBlockServiceNotification(recipientId, data as ISubscriptionNotificationData, queryRunner); await this.notificationsService.createBlockServiceNotification(recipientId, data as ISubscriptionNotificationData, queryRunner);
break; break;
// Payment Notifications // Payment Notifications
case NotifType.PAYMENT_REMINDER: case NotifType.PAYMENT_REMINDER:
await this.notificationsService.createPaymentReminderNotification(recipientId, data as IPaymentNotificationData, queryRunner); await this.notificationsService.createPaymentReminderNotification(recipientId, data as IPaymentNotificationData, queryRunner);
break; break;
case NotifType.PAYMENT_CANCELLATION: case NotifType.PAYMENT_CANCELLATION:
await this.notificationsService.createPaymentCancellationNotification( await this.notificationsService.createPaymentCancellationNotification(recipientId, data as IPaymentNotificationData, queryRunner);
recipientId,
data as IPaymentNotificationData,
queryRunner,
);
break; break;
default: default:
this.logger.warn(`Unknown notification type: ${type}`); this.logger.warn(`Unknown notification type: ${type}`);
} }
await queryRunner.commitTransaction(); await queryRunner.commitTransaction();
clearInterval(heartbeat);
return true; return true;
} catch (error) {
clearInterval(heartbeat);
throw error;
}
} catch (error) { } catch (error) {
this.logger.error( this.logger.error(
`Failed to process notification: ${error instanceof Error ? error.message : "Unknown error"}`, `Failed to process notification: ${error instanceof Error ? error.message : "Unknown error"}`,
@@ -198,6 +167,7 @@ export class NotificationProcessor extends WorkerProcessor {
await queryRunner.rollbackTransaction(); await queryRunner.rollbackTransaction();
throw error; throw error;
} finally { } finally {
if (heartbeat) clearInterval(heartbeat);
await queryRunner.release(); await queryRunner.release();
} }
} }
@@ -87,7 +87,6 @@ export class NotificationQueue {
async addWalletChargeNotification(recipientId: string, data: IWalletNotificationData) { async addWalletChargeNotification(recipientId: string, data: IWalletNotificationData) {
await this.addNotificationJob(NotifType.WALLET_CHARGE, recipientId, data); await this.addNotificationJob(NotifType.WALLET_CHARGE, recipientId, data);
} }
//TODO:USER THIS
async addWalletDeductionNotification(recipientId: string, data: IWalletNotificationData) { async addWalletDeductionNotification(recipientId: string, data: IWalletNotificationData) {
await this.addNotificationJob(NotifType.WALLET_DEDUCTION, recipientId, data); await this.addNotificationJob(NotifType.WALLET_DEDUCTION, recipientId, data);