import { Worker, Job } from "bullmq"; import { redisConnection } from "../connection"; import { SeleniumJobData } from "../queues"; import { io } from "../../socket"; import { runEligibilityProcessor } from "../processors/eligibilityProcessor"; import { runClaimStatusProcessor } from "../processors/claimStatusProcessor"; import { runClaimSubmitProcessor } from "../processors/claimSubmitProcessor"; /** * Emit a job-status event to the socket that enqueued the job (if any). * Falls back to broadcasting when socketId is absent. */ function emitJobUpdate( socketId: string | undefined, jobId: string, jobType: string, status: "active" | "completed" | "failed", payload: Record ) { const event = "job:update"; const data = { jobId, jobType, status, ...payload }; if (socketId && io) { io.to(socketId).emit(event, data); } else if (io) { io.emit(event, data); } } async function processSeleniumJob(job: Job) { const { jobType, userId, socketId, enrichedPayload } = job.data; const jobId = job.id ?? job.name; emitJobUpdate(socketId, jobId, jobType, "active", { message: "Selenium browser starting…", }); try { let result: any; if (jobType === "eligibility-check") { result = await runEligibilityProcessor({ enrichedPayload, userId, insuranceId: job.data.insuranceId!, formFirstName: job.data.formFirstName, formLastName: job.data.formLastName, formDob: job.data.formDob, }); } else if (jobType === "claim-status-check") { result = await runClaimStatusProcessor({ enrichedPayload, insuranceId: job.data.insuranceId!, }); } else if (jobType === "claim-submit") { result = await runClaimSubmitProcessor({ enrichedPayload, files: job.data.files ?? [], claimId: job.data.claimId, variant: "claimsubmit", socketId, }); } else if (jobType === "claim-pre-auth") { result = await runClaimSubmitProcessor({ enrichedPayload, files: job.data.files ?? [], claimId: job.data.claimId, variant: "claim-pre-auth", socketId, }); } else { throw new Error(`Unknown selenium jobType: ${jobType}`); } emitJobUpdate(socketId, jobId, jobType, "completed", { result }); return result; } catch (err: any) { const errorMsg = err?.message ?? String(err); emitJobUpdate(socketId, jobId, jobType, "failed", { error: errorMsg }); throw err; // let BullMQ mark job as failed / retry } } export function startSeleniumWorker() { const worker = new Worker( "selenium-jobs", processSeleniumJob, { connection: redisConnection, concurrency: 1, // mirror the Python semaphore(1) — 1 browser at a time } ); worker.on("completed", (job) => { console.log(`[seleniumWorker] job ${job.id} (${job.data.jobType}) completed`); }); worker.on("failed", (job, err) => { console.error( `[seleniumWorker] job ${job?.id} (${job?.data.jobType}) failed:`, err.message ); }); console.log("[seleniumWorker] started"); return worker; }