Files
Zumri-Backend/app/workers/document.worker.js
Sathira Sri Sathara b6b345f245 feat: Implement Phase 2 cross-cutting services with email and notification enhancements
- Refactor email verification and password reset utilities to use new email service.
- Introduce email delivery queue and notification delivery model for better tracking.
- Enhance file validation and storage services for improved security and ownership management.
- Add cron job for cleaning inactive notifications with retention policy.
- Update document worker to handle document generation and storage more efficiently.
- Implement logging improvements in activity and log workers.
- Create comprehensive documentation for new API endpoints and services.
- Add unit tests for file validation and notification policies to ensure robustness.
2026-09-03 14:22:34 +05:30

27 lines
2.2 KiB
JavaScript

const { Worker } = require("bullmq");
const crypto = require("crypto");
const connection = require("../config/redisClient");
const db = require("../models");
const { generateDocument } = require("../logic/documents");
const storage = require("../services/storage/storage.service");
module.exports = function createDocumentWorker() {
const worker = new Worker("document-generation", async (job) => {
const { documentId, document, documentType, data } = job.data;
const record = await db.Document.findByPk(documentId);
if (!record) throw Object.assign(new Error("Document record not found"), { code: "DOCUMENT_NOT_FOUND" });
if (record.status === "COMPLETED") return { documentId };
await record.update({ status: "PROCESSING", failure_code: null, failed_at: null });
const { fileBuffer, fileName, mimeType } = await generateDocument(document, documentType, data);
const checksum = crypto.createHash("sha256").update(fileBuffer).digest("hex");
const objectKey = storage.createObjectKey({ ownerId: record.owner_id, mimeType, purpose: "document" });
await storage.uploadBuffer({ buffer: fileBuffer, objectKey, mimeType, checksum });
let upload;
try { upload = await db.Upload.create({ file_path: objectKey, file_type: mimeType, file_size: fileBuffer.length, original_name: fileName, safe_name: fileName, checksum, use_for: "document", uploaded_by: record.created_by, owner_type: record.owner_type, owner_id: record.owner_id, visibility: "PRIVATE", status: "AVAILABLE" }); await record.update({ status: "COMPLETED", storage_upload_id: upload.id, generated_at: new Date() }); }
catch (error) { await storage.deleteObject(objectKey).catch(() => undefined); throw error; }
return { documentId, uploadId: upload.id };
}, { connection, concurrency: Number(process.env.DOCUMENT_QUEUE_CONCURRENCY || 2) });
worker.on("failed", async (job, error) => { if (job && job.attemptsMade >= (job.opts.attempts || 1)) await db.Document.update({ status: "FAILED", failed_at: new Date(), failure_code: String(error.code || error.name || "GENERATION_FAILED").slice(0, 100) }, { where: { doc_id: job.data.documentId } }).catch(() => undefined); console.error("Document generation failed", { jobId: job?.id, code: error.code || error.name }); });
return worker;
};