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 { 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 { try { await heartbeat(workerId); } catch (error) { await recoverRedisConnection( error instanceof Error ? error.message : "heartbeat", ); await heartbeat(workerId); } } async function bootstrap(): Promise { 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); });