first
This commit is contained in:
@@ -0,0 +1,176 @@
|
||||
import { Job, Worker } from "bullmq";
|
||||
import { startSession } from "mongoose";
|
||||
|
||||
import { OrderItemsStatus, OrdersStatus, PaymentStatus } from "../../common/enums/order.enum";
|
||||
import { Logger } from "../../core/logging/logger";
|
||||
import { connectRedis } from "../../db/connection";
|
||||
import { OrderModel } from "../../modules/order/models/order.model";
|
||||
import { OrderItemModel } from "../../modules/order/models/orderItem.model";
|
||||
import { CartPaymentModel } from "../../modules/payment/models/payments.model";
|
||||
import { ProductVariantModel } from "../../modules/product/models/productVariant.model";
|
||||
import { IUser } from "../../modules/user/models/Abstraction/IUser";
|
||||
import { SMS } from "../../utils/sms.service";
|
||||
import { QUEUE_KEYS } from "../constant";
|
||||
import { OrderQueue } from "./OrderQueue";
|
||||
|
||||
class OrderProcessor {
|
||||
private static instance: OrderProcessor;
|
||||
public worker: Worker;
|
||||
private logger = new Logger();
|
||||
private constructor() {}
|
||||
|
||||
public static getInstance(): OrderProcessor {
|
||||
if (!OrderProcessor.instance) {
|
||||
OrderProcessor.instance = new OrderProcessor();
|
||||
}
|
||||
return OrderProcessor.instance;
|
||||
}
|
||||
|
||||
async processQueue() {
|
||||
this.worker = new Worker(
|
||||
QUEUE_KEYS.QUEUE_NAME.ORDER_QUEUE,
|
||||
async (job: Job) => {
|
||||
const { orderId } = job.data;
|
||||
switch (job.name) {
|
||||
case QUEUE_KEYS.JOBS.CHECK_ORDER_STATUS:
|
||||
await this.checkAndNotify(orderId);
|
||||
break;
|
||||
case QUEUE_KEYS.JOBS.CANCEL_ORDER:
|
||||
await this.cancelPendingOrder(orderId);
|
||||
break;
|
||||
case QUEUE_KEYS.JOBS.SET_MAIN_STATUS:
|
||||
await this.setMainStatus(orderId);
|
||||
break;
|
||||
}
|
||||
},
|
||||
{ connection: connectRedis() },
|
||||
);
|
||||
|
||||
this.worker.on("completed", (job) => {
|
||||
this.logger.info(`${job.id} has completed!`);
|
||||
});
|
||||
|
||||
this.worker.on("failed", (job, err) => {
|
||||
this.logger.error(`${job?.id} has failed with ${err.message}`, err);
|
||||
});
|
||||
|
||||
this.worker.on("error", (err) => {
|
||||
this.logger.error(`Worker error: ${err.message}`, err);
|
||||
});
|
||||
}
|
||||
|
||||
private async setMainStatus(orderId: number) {
|
||||
try {
|
||||
this.logger.info(`Checking status for order ID: ${orderId}`);
|
||||
const order = await OrderModel.findById(orderId);
|
||||
if (!order) throw new Error("Order not found");
|
||||
|
||||
const orderItems = await OrderItemModel.find({ order: order._id });
|
||||
const allDelivered = orderItems.every((item) => item.status === OrderItemsStatus.Delivered);
|
||||
const allCancelled = orderItems.every((item) =>
|
||||
[OrderItemsStatus.cancelled_shop, OrderItemsStatus.cancelled_user, OrderItemsStatus.cancelled_system].includes(item.status),
|
||||
);
|
||||
|
||||
if (allDelivered) {
|
||||
await OrderModel.findByIdAndUpdate(orderId, { orderStatus: OrdersStatus.Delivered });
|
||||
} else if (allCancelled) {
|
||||
await OrderModel.findByIdAndUpdate(orderId, { orderStatus: OrdersStatus.Cancelled });
|
||||
}
|
||||
} catch (error) {
|
||||
this.logger.error("Error updating order status:", error);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
private async checkAndNotify(orderId: number) {
|
||||
try {
|
||||
this.logger.info(`Checking status for order ID: ${orderId}`);
|
||||
const order = await OrderModel.findById(orderId);
|
||||
if (!order) {
|
||||
this.logger.warn(`Order ${orderId} not found in database`);
|
||||
throw new Error("Order not found");
|
||||
}
|
||||
|
||||
const user = order.user as unknown as IUser;
|
||||
this.logger.info("User phone is:", user.phoneNumber);
|
||||
|
||||
if (order.orderStatus === OrdersStatus.wait_payment) {
|
||||
await SMS.sendOrderPendingSms(user.phoneNumber ?? order.shipmentAddress.phone, user.fullName, order._id);
|
||||
this.logger.info(`Order ${orderId} is still pending. SMS sent to user.`);
|
||||
await OrderQueue.addOrderToCancelQueue(orderId);
|
||||
}
|
||||
} catch (error) {
|
||||
this.logger.error(`Error in checkAndNotify for order ${orderId}:`, error);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
private async cancelPendingOrder(orderId: number) {
|
||||
const session = await startSession();
|
||||
session.startTransaction();
|
||||
try {
|
||||
this.logger.info(`Attempting to cancel order ID: ${orderId}`);
|
||||
const order = await OrderModel.findById(orderId).session(session);
|
||||
if (!order) {
|
||||
this.logger.info(`Order ${orderId} not found.`);
|
||||
throw new Error("Order not found");
|
||||
}
|
||||
|
||||
const payment = await CartPaymentModel.findById(order.payment).session(session);
|
||||
if (!payment) {
|
||||
this.logger.info(`Payment of order ${orderId} not found.`);
|
||||
throw new Error("Payment not found");
|
||||
}
|
||||
|
||||
const orderItems = await OrderItemModel.find({ order: order._id });
|
||||
|
||||
if (order.orderStatus === OrdersStatus.wait_payment) {
|
||||
order.orderStatus = OrdersStatus.cancelled_system;
|
||||
payment.paymentStatus = PaymentStatus.Cancelled;
|
||||
|
||||
await OrderItemModel.updateMany({ order: order._id }, { status: OrderItemsStatus.cancelled_system }, { session });
|
||||
|
||||
const items = orderItems.flatMap((sellerItem) =>
|
||||
sellerItem.shipmentItems.map((shipmentItem) => ({
|
||||
variantId: shipmentItem.variant.toString(),
|
||||
quantity: shipmentItem.quantity,
|
||||
})),
|
||||
);
|
||||
|
||||
const bulkOperations = items.map(({ variantId, quantity }) => ({
|
||||
updateOne: {
|
||||
filter: { _id: variantId },
|
||||
update: { $inc: { stock: quantity } },
|
||||
},
|
||||
}));
|
||||
|
||||
if (bulkOperations.length > 0) {
|
||||
await ProductVariantModel.bulkWrite(bulkOperations, { session });
|
||||
}
|
||||
|
||||
await order.save({ session });
|
||||
await payment.save({ session });
|
||||
|
||||
this.logger.info(`Order ${orderId} cancelled due to no payment after 20 minutes.`);
|
||||
}
|
||||
|
||||
await session.commitTransaction();
|
||||
} catch (error) {
|
||||
await session.abortTransaction();
|
||||
this.logger.error(`Error in cancelPendingOrder for order ${orderId}:`, error);
|
||||
throw error;
|
||||
} finally {
|
||||
await session.endSession();
|
||||
}
|
||||
}
|
||||
|
||||
async closeWorker() {
|
||||
if (this.worker) {
|
||||
await this.worker.close();
|
||||
this.logger.warn("Worker closed successfully.");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
const instance = OrderProcessor.getInstance();
|
||||
export { instance as OrderProcessor };
|
||||
@@ -0,0 +1,50 @@
|
||||
import { JobsOptions, Queue } from "bullmq";
|
||||
|
||||
import { connectRedis } from "../../db/connection";
|
||||
import { QUEUE_KEYS } from "../constant";
|
||||
|
||||
class OrderQueue {
|
||||
private DEFAULT_REMOVE_CONFIG: JobsOptions = {
|
||||
removeOnComplete: {
|
||||
age: 3600,
|
||||
},
|
||||
removeOnFail: {
|
||||
age: 2 * 3600,
|
||||
},
|
||||
// attempts: 3,
|
||||
// backoff: {
|
||||
// type: "exponential",
|
||||
// delay: 2 * 60 * 1000,
|
||||
// },
|
||||
};
|
||||
private static instance: OrderQueue;
|
||||
private queue: Queue;
|
||||
|
||||
private constructor() {
|
||||
this.queue = new Queue(QUEUE_KEYS.QUEUE_NAME.ORDER_QUEUE, {
|
||||
connection: connectRedis(),
|
||||
});
|
||||
}
|
||||
|
||||
public static getInstance(): OrderQueue {
|
||||
if (!OrderQueue.instance) {
|
||||
OrderQueue.instance = new OrderQueue();
|
||||
}
|
||||
return OrderQueue.instance;
|
||||
}
|
||||
|
||||
async addOrderToQueue(orderId: number, delay: number = 10 * 60 * 1000) {
|
||||
await this.queue.add(QUEUE_KEYS.JOBS.CHECK_ORDER_STATUS, { orderId }, { delay, ...this.DEFAULT_REMOVE_CONFIG });
|
||||
}
|
||||
|
||||
async addOrderToCheckStatus(orderId: number, delay: number = 5 * 1000) {
|
||||
await this.queue.add(QUEUE_KEYS.JOBS.SET_MAIN_STATUS, { orderId }, { delay, ...this.DEFAULT_REMOVE_CONFIG });
|
||||
}
|
||||
|
||||
async addOrderToCancelQueue(orderId: number, delay: number = 10 * 60 * 1000) {
|
||||
await this.queue.add(QUEUE_KEYS.JOBS.CANCEL_ORDER, { orderId }, { delay, ...this.DEFAULT_REMOVE_CONFIG });
|
||||
}
|
||||
}
|
||||
|
||||
const instance = OrderQueue.getInstance();
|
||||
export { instance as OrderQueue };
|
||||
Reference in New Issue
Block a user