first commit
This commit is contained in:
@@ -0,0 +1,40 @@
|
||||
/**
|
||||
* 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/activity.worker.js
|
||||
|
||||
const { Worker } = require("bullmq");
|
||||
const connection = require("../config/redisClient");
|
||||
const db = require("../models");
|
||||
|
||||
const UserActivity = db.UserActivity;
|
||||
|
||||
const createActivityWorker = () => {
|
||||
const worker = new Worker(
|
||||
"activity-queue",
|
||||
async (job) => {
|
||||
await UserActivity.create(job.data);
|
||||
},
|
||||
{
|
||||
connection,
|
||||
},
|
||||
);
|
||||
|
||||
worker.on("completed", (job) => {
|
||||
console.log(`✅ Activity logged: ${job.id}`);
|
||||
});
|
||||
|
||||
worker.on("failed", (job, err) => {
|
||||
console.error(`❌ Activity failed: ${job.id}`, err.message);
|
||||
});
|
||||
|
||||
return worker;
|
||||
};
|
||||
|
||||
module.exports = createActivityWorker;
|
||||
@@ -0,0 +1,96 @@
|
||||
/**
|
||||
* 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 { generateDocument } = require("../logic/documents");
|
||||
|
||||
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;
|
||||
};
|
||||
@@ -0,0 +1,34 @@
|
||||
/**
|
||||
* 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/index.js
|
||||
|
||||
require("dotenv").config()
|
||||
const createActivityWorker = require("./activity.worker");
|
||||
const createLogWorker = require("./log.worker");
|
||||
const createDocumentWorker = require("./document.worker");
|
||||
|
||||
// Import future workers here
|
||||
// const createEmailWorker = require("./email.worker");
|
||||
|
||||
console.log("🚀 Starting workers...");
|
||||
|
||||
// Initialize workers (async now)
|
||||
(async () => {
|
||||
try {
|
||||
createActivityWorker();
|
||||
createLogWorker();
|
||||
await createDocumentWorker();
|
||||
|
||||
console.log("✅ All workers started");
|
||||
} catch (err) {
|
||||
console.error("❌ Error starting workers:", err.message);
|
||||
process.exit(1);
|
||||
}
|
||||
})();
|
||||
@@ -0,0 +1,73 @@
|
||||
/**
|
||||
* 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/log.worker.js
|
||||
|
||||
const { Worker } = require("bullmq");
|
||||
const fs = require("fs");
|
||||
const path = require("path");
|
||||
|
||||
const connection = require("../config/redisClient");
|
||||
|
||||
const logsDirectory = path.join(__dirname, "../../logs");
|
||||
|
||||
// Ensure logs directory exists
|
||||
fs.mkdirSync(logsDirectory, { recursive: true });
|
||||
|
||||
const createLogWorker = () => {
|
||||
const worker = new Worker(
|
||||
"logQueue",
|
||||
|
||||
async (job) => {
|
||||
const { message, timestamp } = job.data;
|
||||
|
||||
const logDate = new Date(timestamp)
|
||||
.toISOString()
|
||||
.split("T")[0];
|
||||
|
||||
const logFilePath = path.join(
|
||||
logsDirectory,
|
||||
`app-${logDate}.log`
|
||||
);
|
||||
|
||||
const logLine = `[${timestamp}] ${message}\n`;
|
||||
|
||||
await fs.promises.appendFile(
|
||||
logFilePath,
|
||||
logLine,
|
||||
"utf8"
|
||||
);
|
||||
},
|
||||
|
||||
{
|
||||
connection,
|
||||
}
|
||||
);
|
||||
|
||||
worker.on("completed", (job) => {
|
||||
console.log(`Log job ${job.id} completed`);
|
||||
});
|
||||
|
||||
worker.on("failed", (job, err) => {
|
||||
console.error(
|
||||
`Log job ${job?.id} failed:`,
|
||||
err
|
||||
);
|
||||
});
|
||||
|
||||
worker.on("error", (err) => {
|
||||
console.error("Log worker error:", err);
|
||||
});
|
||||
|
||||
console.log("🟢 Log worker started");
|
||||
|
||||
return worker;
|
||||
};
|
||||
|
||||
module.exports = createLogWorker;
|
||||
Reference in New Issue
Block a user