You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
244 lines
8.1 KiB
244 lines
8.1 KiB
import { loadDevEnv } from "@/lib/dev-env";
|
|
|
|
loadDevEnv();
|
|
import { assertRpaWorkerEnv } from "@/lib/rpa/env";
|
|
import { ensureNativeInfraIfNeeded } from "@/lib/infra/ensure-native-infra";
|
|
import {
|
|
ensureRedisReady,
|
|
getRedisConnectionOptions,
|
|
resetRedisClient,
|
|
} from "@/lib/redis";
|
|
import { Worker } from "bullmq";
|
|
import os from "node:os";
|
|
import {
|
|
RPA_QUEUE_NAME,
|
|
RPA_WORKER_CONCURRENCY,
|
|
} from "@/lib/constants/rpa";
|
|
import { isCircuitOpen } from "@/modules/cache/circuit-breaker";
|
|
import type { AddressCandidatesJobData } from "@/modules/address/candidates-queue-types";
|
|
import type { ParkedSessionJobData } from "@/modules/address/parked-session-queue-types";
|
|
import { processAddressCandidatesJob } from "@/workers/rpa/candidates-job-handler";
|
|
import { processQuoteJob } from "@/workers/rpa/job-handler";
|
|
import type { Priority1QuoteJobData } from "@/modules/priority1/rpa-queue";
|
|
import { processPriority1QuoteJob } from "@/workers/rpa/priority1-job-handler";
|
|
import {
|
|
processParkedSessionPreheatJob,
|
|
processReleaseParkedSessionJob,
|
|
} from "@/workers/rpa/parked-session-job-handler";
|
|
import { processSessionConfirmJob } from "@/workers/rpa/session-confirm-job-handler";
|
|
import type { SessionConfirmJobData } from "@/modules/address/session-confirm-queue-types";
|
|
import type { QuoteJobData } from "@/modules/quote/rpa-queue";
|
|
import { closeBrowser } from "@/workers/rpa/session-manager";
|
|
import { sweepExpiredParkedSessions } from "@/workers/rpa/parked-quote-session";
|
|
import { isParkedSessionEnabled } from "@/lib/rpa/env";
|
|
import { getRpaAddressMode } from "@/lib/rpa/address-mode";
|
|
import {
|
|
getRpaWorkerLockDurationMs,
|
|
} from "@/workers/rpa/queue";
|
|
import { recoverOrphanActiveJobs } from "@/workers/rpa/queue-recovery";
|
|
import {
|
|
heartbeat,
|
|
isWorkerPaused,
|
|
} from "@/workers/rpa/worker-state";
|
|
|
|
const WORKER_ID =
|
|
process.env.RPA_WORKER_ID?.trim() || `${os.hostname()}-${process.pid}`;
|
|
process.env.RPA_WORKER_ID = WORKER_ID;
|
|
|
|
const HEARTBEAT_INTERVAL_MS = 10_000;
|
|
const CIRCUIT_CHECK_INTERVAL_MS = 5_000;
|
|
const NATIVE_INFRA_SETTLE_MS = 3_000;
|
|
const PARKED_SWEEP_INTERVAL_MS = 60_000;
|
|
|
|
async function recoverRedisConnection(reason: string): Promise<void> {
|
|
console.warn(`[rpa-worker] Redis 恢复中(${reason})…`);
|
|
try {
|
|
ensureNativeInfraIfNeeded();
|
|
} catch (error) {
|
|
console.error("[rpa-worker] WSL 基础设施恢复失败:", error);
|
|
}
|
|
resetRedisClient();
|
|
await ensureRedisReady(15, 2_000);
|
|
}
|
|
|
|
async function safeHeartbeat(workerId: string): Promise<void> {
|
|
try {
|
|
await heartbeat(workerId);
|
|
} catch (error) {
|
|
await recoverRedisConnection(
|
|
error instanceof Error ? error.message : "heartbeat",
|
|
);
|
|
await heartbeat(workerId);
|
|
}
|
|
}
|
|
|
|
async function bootstrap(): Promise<void> {
|
|
assertRpaWorkerEnv();
|
|
ensureNativeInfraIfNeeded();
|
|
if (
|
|
process.platform === "win32" &&
|
|
process.env.DEV_INFRA_MODE?.trim().toLowerCase() === "native"
|
|
) {
|
|
await new Promise((resolve) => setTimeout(resolve, NATIVE_INFRA_SETTLE_MS));
|
|
}
|
|
await ensureRedisReady();
|
|
const recovered = await recoverOrphanActiveJobs();
|
|
if (recovered > 0) {
|
|
console.warn(`[rpa-worker] 已回收 ${recovered} 条 orphan active 队列 job`);
|
|
}
|
|
console.log(`[rpa-worker] 启动 worker_id=${WORKER_ID}`);
|
|
console.log(`[rpa-worker] RPA_ADDRESS_MODE=${getRpaAddressMode()}`);
|
|
console.log(
|
|
`[rpa-worker] 驻留页加速=${isParkedSessionEnabled() ? "已启用(RPA_PARKED_SESSION=true)" : "已关闭(默认)"}`,
|
|
);
|
|
|
|
if (process.env.RPA_MOCK_MODE !== "true") {
|
|
try {
|
|
const { ensureMothershipStorageState } = await import(
|
|
"@/lib/rpa/ensure-storage-state"
|
|
);
|
|
await ensureMothershipStorageState();
|
|
} catch (error) {
|
|
const brief = error instanceof Error ? error.message : String(error);
|
|
console.warn(
|
|
`[rpa-worker] 启动时自动获取 MotherShip 会话失败,首单将重试: ${brief.slice(0, 160)}`,
|
|
);
|
|
}
|
|
}
|
|
|
|
const worker = new Worker(
|
|
RPA_QUEUE_NAME,
|
|
async (job) => {
|
|
if (job.name === "address-candidates") {
|
|
console.log(
|
|
`[rpa-worker] 消费 address-candidates job_id=${(job.data as AddressCandidatesJobData).jobId}`,
|
|
);
|
|
await processAddressCandidatesJob(
|
|
job.data as AddressCandidatesJobData,
|
|
);
|
|
return;
|
|
}
|
|
|
|
if (job.name === "parked-session-preheat") {
|
|
await processParkedSessionPreheatJob(job.data as ParkedSessionJobData);
|
|
return;
|
|
}
|
|
|
|
if (job.name === "release-parked-session") {
|
|
await processReleaseParkedSessionJob(job.data as ParkedSessionJobData);
|
|
return;
|
|
}
|
|
|
|
if (job.name === "session-confirm-addresses") {
|
|
await processSessionConfirmJob(job.data as SessionConfirmJobData);
|
|
return;
|
|
}
|
|
|
|
if (job.name === "priority1-quote") {
|
|
const p1Data = job.data as Priority1QuoteJobData;
|
|
const maxAttempts = job.opts.attempts ?? 1;
|
|
const isFinalAttempt = job.attemptsMade + 1 >= maxAttempts;
|
|
console.log(
|
|
`[rpa-worker] 消费 priority1-quote quote_id=${p1Data.quoteId} attempt=${job.attemptsMade + 1}/${maxAttempts}`,
|
|
);
|
|
await processPriority1QuoteJob(p1Data, WORKER_ID, { isFinalAttempt });
|
|
return;
|
|
}
|
|
|
|
const quoteData = job.data as QuoteJobData;
|
|
const maxAttempts = job.opts.attempts ?? 1;
|
|
const isFinalAttempt = job.attemptsMade + 1 >= maxAttempts;
|
|
console.log(
|
|
`[rpa-worker] 消费 job quote_id=${quoteData.quoteId} attempt=${job.attemptsMade + 1}/${maxAttempts}`,
|
|
);
|
|
await processQuoteJob(quoteData, WORKER_ID, { isFinalAttempt });
|
|
},
|
|
{
|
|
connection: getRedisConnectionOptions(),
|
|
concurrency: RPA_WORKER_CONCURRENCY,
|
|
lockDuration: getRpaWorkerLockDurationMs(),
|
|
stalledInterval: 30_000,
|
|
maxStalledCount: 2,
|
|
},
|
|
);
|
|
|
|
worker.on("failed", (job, error) => {
|
|
const label =
|
|
job?.name === "address-candidates"
|
|
? `job_id=${(job.data as AddressCandidatesJobData).jobId}`
|
|
: `quote_id=${(job?.data as QuoteJobData | undefined)?.quoteId}`;
|
|
console.error(`[rpa-worker] job 失败 ${label}:`, error.message);
|
|
});
|
|
|
|
worker.on("error", (error) => {
|
|
console.error("[rpa-worker] BullMQ 连接错误:", error.message);
|
|
});
|
|
|
|
worker.on("completed", (job) => {
|
|
if (job.name === "address-candidates") {
|
|
console.log(
|
|
`[rpa-worker] address-candidates 完成 job_id=${(job.data as AddressCandidatesJobData).jobId}`,
|
|
);
|
|
return;
|
|
}
|
|
console.log(
|
|
`[rpa-worker] job 完成 quote_id=${(job.data as QuoteJobData).quoteId}`,
|
|
);
|
|
});
|
|
|
|
await safeHeartbeat(WORKER_ID);
|
|
console.log("[rpa-worker] Redis 心跳就绪");
|
|
|
|
const heartbeatTimer = setInterval(() => {
|
|
safeHeartbeat(WORKER_ID).catch((error) => {
|
|
console.error("[rpa-worker] 心跳失败:", error);
|
|
});
|
|
}, HEARTBEAT_INTERVAL_MS);
|
|
|
|
const circuitTimer = setInterval(async () => {
|
|
try {
|
|
const paused = await isWorkerPaused(WORKER_ID);
|
|
const circuitOpen = await isCircuitOpen();
|
|
if (circuitOpen || paused) {
|
|
if (!worker.isPaused()) {
|
|
await worker.pause(true);
|
|
console.warn(
|
|
`[rpa-worker] 暂停消费 circuit=${circuitOpen} paused=${paused}`,
|
|
);
|
|
}
|
|
} else if (worker.isPaused()) {
|
|
worker.resume();
|
|
console.log("[rpa-worker] 恢复消费");
|
|
}
|
|
} catch (error) {
|
|
console.error("[rpa-worker] 熔断/暂停检查失败:", error);
|
|
await recoverRedisConnection("circuit-check").catch(() => undefined);
|
|
}
|
|
}, CIRCUIT_CHECK_INTERVAL_MS);
|
|
|
|
const parkedSweepTimer = setInterval(() => {
|
|
const removed = sweepExpiredParkedSessions();
|
|
if (removed > 0) {
|
|
console.log(`[rpa-worker] 清理过期驻留会话 ${removed} 条`);
|
|
}
|
|
}, PARKED_SWEEP_INTERVAL_MS);
|
|
|
|
const shutdown = async () => {
|
|
console.log("[rpa-worker] 正在关闭…");
|
|
clearInterval(heartbeatTimer);
|
|
clearInterval(circuitTimer);
|
|
clearInterval(parkedSweepTimer);
|
|
await worker.close();
|
|
await closeBrowser();
|
|
process.exit(0);
|
|
};
|
|
|
|
process.on("SIGINT", shutdown);
|
|
process.on("SIGTERM", shutdown);
|
|
}
|
|
|
|
bootstrap().catch((error) => {
|
|
console.error("[rpa-worker] 启动失败:", error);
|
|
process.exit(1);
|
|
});
|