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.
This commit is contained in:
@@ -1,96 +1,26 @@
|
||||
/**
|
||||
* Copyright (c) 2026 Niolla
|
||||
* All rights reserved.
|
||||
*
|
||||
* This source code is proprietary and confidential.
|
||||
* Unauthorized copying, modification, distribution, or use
|
||||
* of this file, via any medium, is strictly prohibited.
|
||||
*/
|
||||
|
||||
// app/workers/pdf.worker.js
|
||||
|
||||
const { Worker } = require("bullmq");
|
||||
|
||||
const redis = require("../config/redisClient");
|
||||
|
||||
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");
|
||||
|
||||
const { uploadToS3 } = require("../utils/s3Upload.utill");
|
||||
|
||||
module.exports = async function createDocumentWorker() {
|
||||
const documentWorker = new Worker(
|
||||
"document-generation",
|
||||
|
||||
async (job) => {
|
||||
try {
|
||||
const { document, documentType, data } = job.data;
|
||||
|
||||
console.log(
|
||||
`📄 Processing Job ${job.id}: document="${document}", documentType="${documentType}"`,
|
||||
);
|
||||
|
||||
const { fileBuffer, fileName, mimeType } = await generateDocument(
|
||||
document,
|
||||
documentType,
|
||||
data,
|
||||
);
|
||||
|
||||
/**
|
||||
* Store in S3 instead of local filesystem
|
||||
*/
|
||||
const s3Key = await uploadToS3(
|
||||
{
|
||||
originalname: fileName,
|
||||
buffer: fileBuffer,
|
||||
mimetype: mimeType,
|
||||
},
|
||||
"document",
|
||||
);
|
||||
|
||||
// Extract UUID from S3 key
|
||||
const documentId = s3Key
|
||||
.split("/")
|
||||
.pop()
|
||||
.replace(/\.[^/.]+$/, "");
|
||||
|
||||
console.log(`✅ Job ${job.id} completed: ${fileName} (${s3Key})`);
|
||||
|
||||
/**
|
||||
* Keep response compatible with existing flow
|
||||
*/
|
||||
return {
|
||||
fileName,
|
||||
mimeType,
|
||||
size: fileBuffer.length,
|
||||
s3Key,
|
||||
documentId,
|
||||
};
|
||||
} catch (err) {
|
||||
console.error(`❌ Job ${job.id} failed:`, err.message);
|
||||
|
||||
throw err;
|
||||
}
|
||||
},
|
||||
|
||||
{
|
||||
connection: redis,
|
||||
concurrency: 2,
|
||||
},
|
||||
);
|
||||
|
||||
documentWorker.on("error", (err) => {
|
||||
console.error("❌ Document Worker Error:", err);
|
||||
});
|
||||
|
||||
documentWorker.on("failed", (job, err) => {
|
||||
console.error(`❌ Job ${job?.id} failed after retries:`, err.message);
|
||||
});
|
||||
|
||||
documentWorker.on("completed", (job) => {
|
||||
console.log(`✅ Worker completed Job ${job.id}`);
|
||||
});
|
||||
|
||||
console.log("📄 Document Worker initialized");
|
||||
|
||||
return documentWorker;
|
||||
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;
|
||||
};
|
||||
|
||||
Reference in New Issue
Block a user