488 lines
19 KiB
TypeScript
Executable File
488 lines
19 KiB
TypeScript
Executable File
import { InjectQueue, Processor } from "@nestjs/bullmq";
|
|
import { Job, Queue } from "bullmq";
|
|
import dayjs from "dayjs";
|
|
import Decimal from "decimal.js";
|
|
import { DataSource, In, QueryRunner } from "typeorm";
|
|
|
|
import { WorkerProcessor } from "../../../common/queues/worker.processor";
|
|
import { LoggerService } from "../../logger/logger.service";
|
|
import { NotificationQueue } from "../../notifications/queue/notification.queue";
|
|
import { UserSubscription } from "../../subscriptions/entities/user-subscription.entity";
|
|
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";
|
|
import { InvoicePurpose } from "../interfaces/external-invoice.interface";
|
|
@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_RENEWAL_JOB_NAME:
|
|
return this.createSubscriptionRenewalInvoice(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<void> {
|
|
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<void> {
|
|
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<Invoice> {
|
|
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,
|
|
});
|
|
|
|
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<Invoice> {
|
|
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();
|
|
|
|
invoice.currentRecurringCycle += 1;
|
|
await queryRunner.manager.save(Invoice, invoice);
|
|
|
|
return this.invoicesService.createInvoiceForSubscription(invoice.user, plan, userSubscriptionPlan, dueDate, queryRunner,InvoicePurpose.NEW);
|
|
}
|
|
//********************************** */
|
|
|
|
private async createSubscriptionRenewalInvoice(job: Job<{ userSubscriptionId: string }>, token?: string): Promise<void> {
|
|
const { userSubscriptionId } = job.data;
|
|
this.logger.log(`Creating subscription renewal invoice for subscription: ${userSubscriptionId} ${token ? `with token: ${token}` : ""}`);
|
|
|
|
const queryRunner = this.dataSource.createQueryRunner();
|
|
|
|
try {
|
|
await queryRunner.connect();
|
|
await queryRunner.startTransaction();
|
|
|
|
const userSubscription = await queryRunner.manager.findOne(UserSubscription, {
|
|
where: { id: userSubscriptionId },
|
|
relations: {
|
|
user: true,
|
|
plan: { service: true, directDiscount: true },
|
|
},
|
|
});
|
|
|
|
if (!userSubscription) {
|
|
throw new Error(`User subscription not found: ${userSubscriptionId}`);
|
|
}
|
|
|
|
if (userSubscription.status !== SubscriptionStatus.ACTIVE) {
|
|
this.logger.warn(`Subscription ${userSubscriptionId} is not active, skipping renewal invoice creation`);
|
|
await queryRunner.commitTransaction();
|
|
return;
|
|
}
|
|
|
|
const daysUntilExpiry = dayjs(userSubscription.endDate).diff(dayjs(), "day");
|
|
if (daysUntilExpiry > INVOICE.SUBSCRIPTION_RENEWAL_DAYS_BEFORE_EXPIRY) {
|
|
this.logger.warn(
|
|
`Subscription ${userSubscriptionId} is not close to expiry (${daysUntilExpiry} days left), skipping renewal invoice creation`,
|
|
);
|
|
await queryRunner.commitTransaction();
|
|
return;
|
|
}
|
|
|
|
const existingRenewalInvoice = await queryRunner.manager.findOne(Invoice, {
|
|
where: {
|
|
user: { id: userSubscription.user.id },
|
|
items: { subscriptionPlan: { id: userSubscriptionId } },
|
|
status: In([InvoiceStatus.PENDING, InvoiceStatus.WAIT_PAYMENT, InvoiceStatus.OVERDUE]),
|
|
},
|
|
relations: { items: { subscriptionPlan: true } },
|
|
});
|
|
|
|
if (existingRenewalInvoice) {
|
|
this.logger.log(`Renewal invoice already exists for subscription ${userSubscriptionId}, skipping creation`);
|
|
await queryRunner.commitTransaction();
|
|
return;
|
|
}
|
|
|
|
const renewalDueDate = dayjs().add(INVOICE.DUEDATE, "day").toDate();
|
|
const renewalInvoice = await this.invoicesService.createInvoiceForSubscription(
|
|
userSubscription.user,
|
|
userSubscription.plan,
|
|
userSubscription,
|
|
renewalDueDate,
|
|
queryRunner,
|
|
InvoicePurpose.RENEW
|
|
);
|
|
|
|
this.logger.log(`Renewal invoice created successfully for subscription ${userSubscriptionId} with id ${renewalInvoice.id}`);
|
|
|
|
// await this.notificationQueue.addInvoiceCreationNotification(userSubscription.user.id, {
|
|
// invoiceId: renewalInvoice.numericId.toString(),
|
|
// dueDate: renewalInvoice.dueDate,
|
|
// createDate: renewalInvoice.createdAt,
|
|
// price: new Decimal(renewalInvoice.totalPrice).toNumber(),
|
|
// userPhone: userSubscription.user.phone,
|
|
// userEmail: userSubscription.user.email,
|
|
// items: `${userSubscription.plan.service.name} - Renewal`,
|
|
// });
|
|
|
|
await queryRunner.commitTransaction();
|
|
this.logger.log(`Subscription renewal invoice created successfully for subscription ${userSubscriptionId}`);
|
|
} catch (error) {
|
|
this.logger.error(`Failed to create subscription renewal invoice for ${userSubscriptionId}:`, error);
|
|
await queryRunner.rollbackTransaction();
|
|
throw error;
|
|
} finally {
|
|
await queryRunner.release();
|
|
}
|
|
}
|
|
//********************************** */
|
|
|
|
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.CANCELED;
|
|
|
|
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<void> {
|
|
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, undefined, queryRunner);
|
|
}
|
|
}
|