feat: add BullMQ queue infrastructure and frontend job status hook
- apps/Backend/src/queue/: connection, queues, workers, processors - apps/Frontend/src/hooks/use-job-status.ts: WebSocket job progress hook - apps/Frontend/src/lib/socket.ts: shared Socket.IO singleton Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,55 @@
|
||||
import { Worker, Job } from "bullmq";
|
||||
import { redisConnection } from "../connection";
|
||||
import { OcrJobData } from "../queues";
|
||||
import { io } from "../../socket";
|
||||
import { runOcrProcessor } from "../processors/ocrProcessor";
|
||||
|
||||
function emitJobUpdate(
|
||||
socketId: string | undefined,
|
||||
jobId: string,
|
||||
status: "active" | "completed" | "failed",
|
||||
payload: Record<string, any>
|
||||
) {
|
||||
const event = "job:update";
|
||||
const data = { jobId, jobType: "ocr", status, ...payload };
|
||||
if (socketId && io) {
|
||||
io.to(socketId).emit(event, data);
|
||||
} else if (io) {
|
||||
io.emit(event, data);
|
||||
}
|
||||
}
|
||||
|
||||
async function processOcrJob(job: Job<OcrJobData>) {
|
||||
const { socketId, files } = job.data;
|
||||
const jobId = job.id ?? job.name;
|
||||
|
||||
emitJobUpdate(socketId, jobId, "active", { message: "OCR processing started…" });
|
||||
|
||||
try {
|
||||
const rows = await runOcrProcessor({ files });
|
||||
emitJobUpdate(socketId, jobId, "completed", { result: { rows } });
|
||||
return rows;
|
||||
} catch (err: any) {
|
||||
const errorMsg = err?.message ?? String(err);
|
||||
emitJobUpdate(socketId, jobId, "failed", { error: errorMsg });
|
||||
throw err;
|
||||
}
|
||||
}
|
||||
|
||||
export function startOcrWorker() {
|
||||
const worker = new Worker<OcrJobData>("ocr-jobs", processOcrJob, {
|
||||
connection: redisConnection,
|
||||
concurrency: 2, // OCR service allows 2 concurrent
|
||||
});
|
||||
|
||||
worker.on("completed", (job) => {
|
||||
console.log(`[ocrWorker] job ${job.id} completed`);
|
||||
});
|
||||
|
||||
worker.on("failed", (job, err) => {
|
||||
console.error(`[ocrWorker] job ${job?.id} failed:`, err.message);
|
||||
});
|
||||
|
||||
console.log("[ocrWorker] started");
|
||||
return worker;
|
||||
}
|
||||
@@ -0,0 +1,104 @@
|
||||
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<string, any>
|
||||
) {
|
||||
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<SeleniumJobData>) {
|
||||
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",
|
||||
});
|
||||
} else if (jobType === "claim-pre-auth") {
|
||||
result = await runClaimSubmitProcessor({
|
||||
enrichedPayload,
|
||||
files: job.data.files ?? [],
|
||||
claimId: job.data.claimId,
|
||||
variant: "claim-pre-auth",
|
||||
});
|
||||
} 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<SeleniumJobData>(
|
||||
"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;
|
||||
}
|
||||
Reference in New Issue
Block a user