chore: add all notification sewrvice to the queue

This commit is contained in:
mahyargdz
2025-05-04 16:06:14 +03:30
parent d50b18b8aa
commit c3ad2a1733
16 changed files with 539 additions and 277 deletions
+38 -57
View File
@@ -6,7 +6,7 @@ import { DataSource, QueryRunner } from "typeorm";
import { WorkerProcessor } from "../../../common/queues/worker.processor";
import { LoggerService } from "../../logger/logger.service";
import { NotificationsService } from "../../notifications/providers/notifications.service";
import { NotificationQueue } from "../../notifications/queue/notification.queue";
import { SubscriptionStatus } from "../../subscriptions/enums/subscription-status.enum";
import { User } from "../../users/entities/user.entity";
import { RoleEnum } from "../../users/enums/role.enum";
@@ -14,14 +14,13 @@ import { INVOICE } from "../constants";
import { Invoice } from "../entities/invoice.entity";
import { InvoiceStatus } from "../enums/invoice-status.enum";
import { InvoicesService } from "../providers/invoices.service";
@Processor(INVOICE.QUEUE_NAME, { concurrency: 5 })
@Processor(INVOICE.QUEUE_NAME, { concurrency: 2 })
export class InvoiceProcessor extends WorkerProcessor {
constructor(
@InjectQueue(INVOICE.QUEUE_NAME) private readonly invoiceQueue: Queue,
private readonly invoicesService: InvoicesService,
private readonly dataSource: DataSource,
private readonly notificationService: NotificationsService,
private readonly notificationQueue: NotificationQueue,
protected readonly logger: LoggerService,
) {
super();
@@ -161,7 +160,7 @@ export class InvoiceProcessor extends WorkerProcessor {
await this.applyLateFeeIfOverdue(invoice, queryRunner);
await this.sendReminderNotification(invoice, queryRunner);
await this.sendReminderNotification(invoice);
await this.scheduleNextReminder(invoice.id);
@@ -204,16 +203,12 @@ export class InvoiceProcessor extends WorkerProcessor {
await queryRunner.manager.save(userSubscription);
await this.notificationService.createBlockServiceNotification(
invoice.user.id,
{
userPhone: invoice.user.phone,
userEmail: invoice.user.email,
planName: userSubscription.plan.name,
invoiceId: invoice.numericId.toString(),
},
queryRunner,
);
await this.notificationQueue.addBlockServiceNotification(invoice.user.id, {
userPhone: invoice.user.phone,
userEmail: invoice.user.email,
planName: userSubscription.plan.name,
invoiceId: invoice.numericId.toString(),
});
}
this.logger.log(`Invoice ${invoice.id} has been cancelled as it is more than 10 days overdue.`);
}
@@ -253,40 +248,32 @@ export class InvoiceProcessor extends WorkerProcessor {
await queryRunner.manager.save(Invoice, invoice);
// Send overdue notification to the user
await this.notificationService.createInvoiceOverdueNotification(
invoice.user.id,
{
userPhone: invoice.user.phone,
userEmail: invoice.user.email,
invoiceId: invoice.numericId.toString(),
price: new Decimal(invoice.totalPrice).toNumber(),
lateFee: new Decimal(newTotalFine).toNumber(),
dueDate: invoice.dueDate,
createDate: invoice.createdAt,
items: invoice.items.map((item) => item.name).join(", "),
},
queryRunner,
);
await this.notificationQueue.addInvoiceOverdueNotification(invoice.user.id, {
userPhone: invoice.user.phone,
userEmail: invoice.user.email,
invoiceId: invoice.numericId.toString(),
price: new Decimal(invoice.totalPrice).toNumber(),
lateFee: new Decimal(newTotalFine).toNumber(),
dueDate: invoice.dueDate,
createDate: invoice.createdAt,
items: invoice.items.map((item) => item.name).join(", "),
});
}
}
}
//********************************** */
private async sendReminderNotification(invoice: Invoice, queryRunner: QueryRunner) {
private async sendReminderNotification(invoice: Invoice) {
this.logger.log(`Sending invoice reminder to user: ${invoice.user.id}`);
await this.notificationService.createBillInvoiceReminderNotification(
invoice.user.id,
{
userPhone: invoice.user.phone,
userEmail: invoice.user.email,
invoiceId: invoice.numericId.toString(),
price: new Decimal(invoice.totalPrice).toNumber(),
dueDate: invoice.dueDate,
createDate: invoice.createdAt,
items: invoice.items.map((item) => item.name).join(", "),
},
queryRunner,
);
await this.notificationQueue.addBillInvoiceReminderNotification(invoice.user.id, {
userPhone: invoice.user.phone,
userEmail: invoice.user.email,
invoiceId: invoice.numericId.toString(),
price: new Decimal(invoice.totalPrice).toNumber(),
dueDate: invoice.dueDate,
createDate: invoice.createdAt,
items: invoice.items.map((item) => item.name).join(", "),
});
}
//********************************** */
@@ -324,9 +311,7 @@ export class InvoiceProcessor extends WorkerProcessor {
items: invoice.items.map((item) => item.name).join(", "),
};
await Promise.all(
admins.map((admin) => this.notificationService.createRecurringInvoiceNotification(admin.id, notificationPayload, queryRunner)),
);
await Promise.all(admins.map((admin) => this.notificationQueue.addRecurringInvoiceNotification(admin.id, notificationPayload)));
}
//********************************** */
@@ -344,17 +329,13 @@ export class InvoiceProcessor extends WorkerProcessor {
const superAdmins = await this.fetchAdmin(queryRunner);
for (const admin of superAdmins) {
await this.notificationService.createNewSubscriptionNotification(
admin.id,
{
date: invoice.dueDate,
userPhone: admin.phone,
userEmail: admin.email,
serviceName: invoice.items[0]?.subscriptionPlan?.plan.service.name ?? "",
fullName: invoice.user.firstName + " " + invoice.user.lastName,
},
queryRunner,
);
await this.notificationQueue.addNewSubscriptionNotification(admin.id, {
date: invoice.dueDate,
userPhone: admin.phone,
userEmail: admin.email,
serviceName: invoice.items[0]?.subscriptionPlan?.plan.service.name ?? "",
fullName: invoice.user.firstName + " " + invoice.user.lastName,
});
}
await queryRunner.commitTransaction();