import { InjectQueue, Processor } from "@nestjs/bullmq"; import { Job, Queue } from "bullmq"; import dayjs from "dayjs"; import Decimal from "decimal.js"; import { DataSource, QueryRunner } from "typeorm"; import { WorkerProcessor } from "../../../common/queues/worker.processor"; import { LoggerService } from "../../logger/logger.service"; import { NotificationQueue } from "../../notifications/queue/notification.queue"; import { SubscriptionStatus } from "../../subscriptions/enums/subscription-status.enum"; import { SupportPlan } from "../../support-plans/entities/support-plan.entity"; import { UserSupportPlan } from "../../support-plans/entities/user-support-plan.entity"; import { UserSupportPlanStatus } from "../../support-plans/enums/user-support-plan-status.enum"; import { User } from "../../users/entities/user.entity"; import { RoleEnum } from "../../users/enums/role.enum"; 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: 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 notificationQueue: NotificationQueue, protected readonly logger: LoggerService, ) { super(); } async process(job: Job, token?: string) { this.logger.log(job.name); switch (job.name) { case INVOICE.REMINDER_JOB_NAME: return this.sendBillInvoiceReminder(job, token); case INVOICE.RECURRING_JOB_NAME: return this.createRecurringInvoice(job, token); case INVOICE.SUBSCRIPTION_ADMIN_NOTIFICATION_JOB_NAME: return this.notifyAdminForSubscriptionInvoice(job); default: this.logger.error(`Unknown job name: ${job.name}`); return; } } //********************************** */ private async createRecurringInvoice(job: Job<{ invoiceId: string; adminCreated: boolean }>, token?: string): Promise { const { invoiceId, adminCreated } = job.data; this.logger.log(`Creating recurring invoice for original invoice: ${invoiceId} ${token ? `with token: ${token}` : ""}`); const queryRunner = this.dataSource.createQueryRunner(); try { await queryRunner.connect(); await queryRunner.startTransaction(); const invoice = await this.fetchInvoice(invoiceId, queryRunner); // eslint-disable-next-line @typescript-eslint/no-unused-expressions adminCreated ? await this.handleAdminCreatedInvoice(invoice, queryRunner) : await this.handleSubscriptionInvoice(invoice, queryRunner); await queryRunner.commitTransaction(); this.logger.log(`Recurring invoice for ${invoiceId} created successfully`); } catch (error) { this.logger.error(`Failed to create recurring invoice for ${invoiceId}:`, error); await queryRunner.rollbackTransaction(); throw error; } finally { await queryRunner.release(); } } //********************************** */ private async validateRecurringInvoice(invoice: Invoice): Promise { if (!invoice.isRecurring) throw new Error(`Invoice ${invoice.id} is not set up for recurring billing`); if (invoice.maxRecurringCycles && invoice.currentRecurringCycle >= invoice.maxRecurringCycles) { throw new Error(`Maximum recurring cycles (${invoice.maxRecurringCycles}) reached for invoice ${invoice.id}`); } if (invoice.status === InvoiceStatus.ARCHIVED) throw new Error(`Cannot create recurring invoice for archived invoice ${invoice.id}`); } //********************************** */ private async handleAdminCreatedInvoice(invoice: Invoice, queryRunner: QueryRunner): Promise { this.logger.verbose(`Creating admin-initiated recurring invoice in draft status`); await this.validateRecurringInvoice(invoice); const newInvoice = queryRunner.manager.create(Invoice, { user: invoice.user, status: InvoiceStatus.DRAFT, dueDate: dayjs().add(INVOICE.DUEDATE, "day").toDate(), totalPrice: invoice.totalPrice, originalPrice: invoice.originalPrice, tax: invoice.tax, items: invoice.items, isRecurring: invoice.isRecurring, recurringPeriod: invoice.recurringPeriod, maxRecurringCycles: invoice.maxRecurringCycles, currentRecurringCycle: invoice.currentRecurringCycle + 1, }); // Update the original invoice's current recurring cycle invoice.currentRecurringCycle += 1; await queryRunner.manager.save(Invoice, invoice); await this.sendNotificationToAdmins(invoice, queryRunner); return queryRunner.manager.save(Invoice, newInvoice); } //********************************** */ private async handleSubscriptionInvoice(invoice: Invoice, queryRunner: QueryRunner): Promise { this.logger.verbose(`Creating subscription-based recurring invoice`); await this.validateRecurringInvoice(invoice); const userSubscriptionPlan = invoice.items[0]?.subscriptionPlan; if (!userSubscriptionPlan) { throw new Error(`No subscription plan found for invoice ${invoice.id}`); } if (userSubscriptionPlan.status !== SubscriptionStatus.ACTIVE) { throw new Error(`Subscription plan ${userSubscriptionPlan.id} is not active`); } const { plan } = userSubscriptionPlan; const dueDate = dayjs().add(INVOICE.DUEDATE, "day").toDate(); // Update the original invoice's current recurring cycle invoice.currentRecurringCycle += 1; await queryRunner.manager.save(Invoice, invoice); return this.invoicesService.createInvoiceForSubscription(invoice.user, plan, userSubscriptionPlan, dueDate, queryRunner); } //********************************** */ private async sendBillInvoiceReminder(job: Job<{ invoiceId: string }>, token?: string) { this.logger.log(`Sending bill invoice reminder: ${job.data.invoiceId} to user. ${token}`); const queryRunner = this.dataSource.createQueryRunner(); try { await queryRunner.connect(); await queryRunner.startTransaction(); const invoice = await this.fetchInvoice(job.data.invoiceId, queryRunner); if (invoice.status === InvoiceStatus.PAID) { this.logger.log(`Invoice ${invoice.id} is already paid.`); await queryRunner.commitTransaction(); return; } if (this.isInvoiceOverdueForCancellation(invoice)) { await this.cancelOverdueInvoice(invoice, queryRunner); await queryRunner.commitTransaction(); return; } await this.applyLateFeeIfOverdue(invoice, queryRunner); await this.sendReminderNotification(invoice); await this.scheduleNextReminder(invoice.id); await queryRunner.commitTransaction(); } catch (error) { this.logger.error(`Failed to process invoice reminder:`, error); await queryRunner.rollbackTransaction(); throw error; } finally { await queryRunner.release(); } } //********************************** */ private async fetchInvoice(invoiceId: string, queryRunner: QueryRunner) { const invoice = await queryRunner.manager.findOne(Invoice, { where: { id: invoiceId }, relations: { user: true, items: { subscriptionPlan: { plan: { service: true } } } }, }); if (!invoice) throw new Error(`Invoice not found: ${invoiceId}`); return invoice; } //********************************** */ private isInvoiceOverdueForCancellation(invoice: Invoice): boolean { const tenDaysAfterDueDate = dayjs(invoice.dueDate).add(INVOICE.MAX_DAYS_AFTER_OVERDUE, "day"); return dayjs().isAfter(tenDaysAfterDueDate); } //********************************** */ private async cancelOverdueInvoice(invoice: Invoice, queryRunner: QueryRunner) { invoice.status = InvoiceStatus.ARCHIVED; await queryRunner.manager.save(Invoice, invoice); if (invoice.items[0]?.subscriptionPlan) { this.logger.log( `Invoice ${invoice.id} has been archived as it is more than 10 days overdue. Subscription Plan ${invoice.items[0].subscriptionPlan.id} has been cancelled.`, ); const userSubscription = invoice.items[0].subscriptionPlan; userSubscription.status = SubscriptionStatus.INACTIVE; await queryRunner.manager.save(userSubscription); await this.notificationQueue.addBlockServiceNotification(invoice.user.id, { userPhone: invoice.user.phone, userEmail: invoice.user.email, planName: userSubscription.plan.name, invoiceId: invoice.numericId.toString(), }); } if (invoice.items[0]?.supportPlan) { this.logger.log( `Invoice ${invoice.id} has been archived as it is more than 10 days overdue. Support Plan ${invoice.items[0].supportPlan.id} has been cancelled.`, ); const userSupportPlan = invoice.items[0].supportPlan; userSupportPlan.status = UserSupportPlanStatus.INACTIVE; await queryRunner.manager.save(UserSupportPlan, userSupportPlan); await this.applyFreeSupportPlan(invoice, queryRunner); } this.logger.log(`Invoice ${invoice.id} has been archived as it is more than 10 days overdue.`); } //********************************** */ private async applyLateFeeIfOverdue(invoice: Invoice, queryRunner: QueryRunner) { const today = dayjs(); const dueDate = dayjs(invoice.dueDate); if (today.isAfter(dueDate)) { // Mark invoice as overdue in status if not already paid if (invoice.status === InvoiceStatus.PENDING || invoice.status === InvoiceStatus.WAIT_PAYMENT) { this.logger.log(`Marking invoice ${invoice.id} as overdue`); invoice.status = InvoiceStatus.OVERDUE; } const daysOverdue = Math.min(today.diff(dueDate, "day"), INVOICE.MAX_DAYS_AFTER_OVERDUE); // Get the original price without late fees to calculate the fine correctly const originalPrice = invoice.originalPrice || invoice.totalPrice; const currentFine = invoice.lateFee || new Decimal(0); const finePerDay = new Decimal(originalPrice).mul(INVOICE.FINE_PERCENTAGE); const newTotalFine = finePerDay.mul(daysOverdue); if (newTotalFine.gt(currentFine)) { this.logger.log(`Applying late fee for invoice ${invoice.id}: ${newTotalFine} (${daysOverdue} days overdue)`); // Calculate the difference between new and current fine const additionalFine = newTotalFine.minus(currentFine); // Update the late fee invoice.lateFee = newTotalFine; // Add only the additional fine to total price to avoid compounding invoice.totalPrice = new Decimal(invoice.totalPrice).plus(additionalFine); await queryRunner.manager.save(Invoice, invoice); // Send overdue notification to the user 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) { this.logger.log(`Sending invoice reminder to user: ${invoice.user.id}`); 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(", "), }); } //********************************** */ private async scheduleNextReminder(invoiceId: string) { this.logger.log(`Invoice ${invoiceId} is still unpaid. Scheduling a new reminder.`); await this.invoiceQueue.add( INVOICE.REMINDER_JOB_NAME, { invoiceId }, { delay: INVOICE.REMINDER_REPEAT_DELAY, // 1 day delay in milliseconds attempts: INVOICE.REMINDER_JOB_ATTEMPTS, backoff: { type: "exponential", delay: INVOICE.REMINDER_JOB_BACKOFF }, }, ); } //********************************** */ private async fetchAdmin(queryRunner: QueryRunner) { return queryRunner.manager.find(User, { where: { roles: { name: RoleEnum.SUPER_ADMIN } }, relations: { roles: true } }); } //********************************** */ private async sendNotificationToAdmins(invoice: Invoice, queryRunner: QueryRunner): Promise { this.logger.log(`Sending notification to admins for invoice ${invoice.id}`); const admins = await this.fetchAdmin(queryRunner); const notificationPayload = { 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(", "), }; await Promise.all(admins.map((admin) => this.notificationQueue.addRecurringInvoiceNotification(admin.id, notificationPayload))); } //********************************** */ private async notifyAdminForSubscriptionInvoice(job: Job<{ invoiceId: string }>) { const queryRunner = this.dataSource.createQueryRunner(); try { await queryRunner.connect(); await queryRunner.startTransaction(); const { invoiceId } = job.data; const invoice = await this.fetchInvoice(invoiceId, queryRunner); const superAdmins = await this.fetchAdmin(queryRunner); for (const admin of superAdmins) { await this.notificationQueue.addNewSubscriptionNotification(admin.id, { date: invoice.paidAt ?? invoice.createdAt, userPhone: admin.phone, userEmail: admin.email, serviceName: invoice.items[0]?.subscriptionPlan?.plan.service.name ?? "", fullName: invoice.user.firstName + " " + invoice.user.lastName, }); } await queryRunner.commitTransaction(); return true; } catch (error) { this.logger.error(`Failed to notify admins for subscription invoice:`, error); await queryRunner.rollbackTransaction(); throw error; } finally { await queryRunner.release(); } } //********************************** */ private async applyFreeSupportPlan(invoice: Invoice, queryRunner: QueryRunner) { const freeSupportPlan = await queryRunner.manager.findOne(SupportPlan, { where: { isFree: true }, }); if (!freeSupportPlan) throw new Error(`No free support plan found`); const startDate = dayjs().toDate(); const endDate = dayjs().add(freeSupportPlan.duration, "day").toDate(); const userSupportPlan = queryRunner.manager.create(UserSupportPlan, { user: invoice.user, supportPlan: freeSupportPlan, startDate, endDate, }); await queryRunner.manager.save(UserSupportPlan, userSupportPlan); const invoiceDueDate = dayjs(userSupportPlan.startDate).add(INVOICE.DUEDATE, "day").toDate(); const newInvoice = await this.invoicesService.createInvoiceForSupportPlan( invoice.user, freeSupportPlan, userSupportPlan, invoiceDueDate, queryRunner, ); userSupportPlan.status = UserSupportPlanStatus.ACTIVE; await queryRunner.manager.save(UserSupportPlan, userSupportPlan); await this.invoicesService.payInvoice(newInvoice.id, invoice.user.id, queryRunner); } }