|
|
import { Queue } from "bullmq";
|
|
|
import {
|
|
|
RPA_JOB_ATTEMPTS,
|
|
|
RPA_QUEUE_NAME,
|
|
|
ADDRESS_CANDIDATES_TIMEOUT_MS,
|
|
|
RPA_JOB_TIMEOUT_MS,
|
|
|
} from "@/lib/constants/rpa";
|
|
|
import { getRedisConnectionOptions } from "@/lib/redis";
|
|
|
import type { AddressCandidatesJobData } from "@/modules/address/candidates-queue-types";
|
|
|
import type { SessionConfirmJobData } from "@/modules/address/session-confirm-queue-types";
|
|
|
import type { ParkedSessionJobData } from "@/modules/address/parked-session-queue-types";
|
|
|
import type { FlockQuoteJobData } from "@/modules/flock/rpa-queue";
|
|
|
import type { Priority1QuoteJobData } from "@/modules/priority1/rpa-queue";
|
|
|
import type { QuoteJobData } from "@/modules/quote/rpa-queue";
|
|
|
|
|
|
export type RpaQueueJobData =
|
|
|
| QuoteJobData
|
|
|
| Priority1QuoteJobData
|
|
|
| FlockQuoteJobData
|
|
|
| AddressCandidatesJobData
|
|
|
| ParkedSessionJobData
|
|
|
| SessionConfirmJobData;
|
|
|
|
|
|
let queueInstance: Queue<RpaQueueJobData> | null = null;
|
|
|
|
|
|
export function getQuoteRpaQueue(): Queue<RpaQueueJobData> {
|
|
|
if (!queueInstance) {
|
|
|
queueInstance = new Queue<RpaQueueJobData>(RPA_QUEUE_NAME, {
|
|
|
connection: getRedisConnectionOptions(),
|
|
|
defaultJobOptions: {
|
|
|
attempts: RPA_JOB_ATTEMPTS,
|
|
|
backoff: { type: "fixed", delay: 1_000 },
|
|
|
removeOnComplete: 200,
|
|
|
removeOnFail: 500,
|
|
|
},
|
|
|
});
|
|
|
}
|
|
|
return queueInstance;
|
|
|
}
|
|
|
|
|
|
export async function addQuoteRpaJob(data: QuoteJobData): Promise<void> {
|
|
|
const queue = getQuoteRpaQueue();
|
|
|
await queue.add("quote", data, {
|
|
|
jobId: data.quoteId,
|
|
|
});
|
|
|
}
|
|
|
|
|
|
export async function addPriority1RpaJob(
|
|
|
data: Priority1QuoteJobData,
|
|
|
): Promise<void> {
|
|
|
const queue = getQuoteRpaQueue();
|
|
|
await queue.add("priority1-quote", data, {
|
|
|
jobId: `p1_${data.quoteId}`,
|
|
|
});
|
|
|
}
|
|
|
|
|
|
export async function addFlockRpaJob(data: FlockQuoteJobData): Promise<void> {
|
|
|
const queue = getQuoteRpaQueue();
|
|
|
await queue.add("flock-quote", data, {
|
|
|
jobId: `flock_${data.quoteId}`,
|
|
|
});
|
|
|
}
|
|
|
|
|
|
export async function addAddressCandidatesJob(
|
|
|
data: AddressCandidatesJobData,
|
|
|
): Promise<void> {
|
|
|
const queue = getQuoteRpaQueue();
|
|
|
await queue.add("address-candidates", data, {
|
|
|
jobId: data.jobId,
|
|
|
});
|
|
|
}
|
|
|
|
|
|
export async function addParkedSessionJob(
|
|
|
name: "parked-session-preheat" | "release-parked-session",
|
|
|
sessionId: string,
|
|
|
): Promise<void> {
|
|
|
const queue = getQuoteRpaQueue();
|
|
|
const payload: ParkedSessionJobData = { sessionId };
|
|
|
await queue.add(name, payload, {
|
|
|
// BullMQ 禁止 custom jobId 含 ':'(Redis 6.0.x 下会抛 Custom Id cannot contain :)
|
|
|
jobId: `${name}_${sessionId}`,
|
|
|
removeOnComplete: 50,
|
|
|
removeOnFail: 100,
|
|
|
});
|
|
|
}
|
|
|
|
|
|
export async function addSessionConfirmJob(
|
|
|
data: SessionConfirmJobData,
|
|
|
): Promise<void> {
|
|
|
const queue = getQuoteRpaQueue();
|
|
|
await queue.add("session-confirm-addresses", data, {
|
|
|
jobId: data.jobId,
|
|
|
});
|
|
|
}
|
|
|
|
|
|
export function getRpaWorkerLockDurationMs(): number {
|
|
|
return (
|
|
|
Math.max(RPA_JOB_TIMEOUT_MS, ADDRESS_CANDIDATES_TIMEOUT_MS) + 15_000
|
|
|
);
|
|
|
}
|
|
|
|
|
|
export async function getQueueDepth(): Promise<number> {
|
|
|
const queue = getQuoteRpaQueue();
|
|
|
const counts = await queue.getJobCounts("waiting", "delayed", "active");
|
|
|
return counts.waiting + counts.delayed + counts.active;
|
|
|
}
|
|
|
|
|
|
export type QuoteRpaQueueCounts = {
|
|
|
waiting: number;
|
|
|
active: number;
|
|
|
delayed: number;
|
|
|
failed: number;
|
|
|
completed: number;
|
|
|
};
|
|
|
|
|
|
export async function getQuoteRpaQueueCounts(): Promise<QuoteRpaQueueCounts> {
|
|
|
const queue = getQuoteRpaQueue();
|
|
|
const counts = await queue.getJobCounts(
|
|
|
"waiting",
|
|
|
"active",
|
|
|
"delayed",
|
|
|
"failed",
|
|
|
"completed",
|
|
|
);
|
|
|
return {
|
|
|
waiting: counts.waiting ?? 0,
|
|
|
active: counts.active ?? 0,
|
|
|
delayed: counts.delayed ?? 0,
|
|
|
failed: counts.failed ?? 0,
|
|
|
completed: counts.completed ?? 0,
|
|
|
};
|
|
|
}
|
|
|
|
|
|
export async function closeQuoteRpaQueue(): Promise<void> {
|
|
|
if (queueInstance) {
|
|
|
await queueInstance.close();
|
|
|
queueInstance = null;
|
|
|
}
|
|
|
}
|
|
|
|
|
|
export { RPA_QUEUE_NAME };
|