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.

272 lines
9.4 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 { processFlockQuoteJob } from "@/workers/rpa/flock-job-handler";
import { processPriority1QuoteJob } from "@/workers/rpa/priority1-job-handler";
import type { FlockQuoteJobData } from "@/modules/flock/rpa-queue";
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 { closeFlockWorkerSession } from "@/workers/rpa/flock/worker-session";
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") {
const { isFlockRpaEnabled, isFlockWorkerReuseSession, getFlockQuotesPerEmail, getFlockQuoteMode, hasFlockLoginCredentials } =
await import("@/lib/flock/env");
if (isFlockRpaEnabled()) {
console.log(
`[rpa-worker] Flock quoteMode=${getFlockQuoteMode()} ` +
`hasLogin=${hasFlockLoginCredentials()} ` +
`emailMode=${process.env.FLOCK_EMAIL_MODE ?? "alias"} ` +
`reuseSession=${isFlockWorkerReuseSession()} ` +
`quotesPerEmail=${getFlockQuotesPerEmail()}`,
);
}
}
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;
}
if (job.name === "flock-quote") {
const flockData = job.data as FlockQuoteJobData;
const maxAttempts = job.opts.attempts ?? 1;
const isFinalAttempt = job.attemptsMade + 1 >= maxAttempts;
console.log(
`[rpa-worker] 消费 flock-quote quote_id=${flockData.quoteId} attempt=${job.attemptsMade + 1}/${maxAttempts}`,
);
await processFlockQuoteJob(flockData, 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 closeFlockWorkerSession();
await closeBrowser();
process.exit(0);
};
process.on("SIGINT", shutdown);
process.on("SIGTERM", shutdown);
}
bootstrap().catch((error) => {
console.error("[rpa-worker] 启动失败:", error);
process.exit(1);
});