Files
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

65 lines
2.3 KiB
JavaScript

/**
* Copyright (c) 2026 Niolla
* All rights reserved.
*/
require("dotenv").config();
const { validateEnvironment } = require("../config/env.config");
const startWorkers = async () => {
const env = validateEnvironment();
const { initializeDatabase, closeDatabase } = require("../config/database.lifecycle");
const { initializeRedis, closeRedis } = require("../config/redis.lifecycle");
try {
await initializeDatabase();
await initializeRedis();
} catch (error) {
await Promise.allSettled([closeRedis(), closeDatabase()]);
throw error;
}
const { closeQueues } = require("../config/queue.lifecycle");
const createActivityWorker = require("./activity.worker");
const createLogWorker = require("./log.worker");
const createDocumentWorker = require("./document.worker");
const createEmailWorker = require("./email.worker");
const workers = [createActivityWorker(), createLogWorker(), await createDocumentWorker(), createEmailWorker()];
console.log("All workers started.");
let shuttingDown = false;
const shutdown = async (reason, exitCode = 0) => {
if (shuttingDown) return;
shuttingDown = true;
console.log(`Workers shutting down (${reason}).`);
const forceTimer = setTimeout(() => process.exit(1), env.SHUTDOWN_TIMEOUT_MS);
forceTimer.unref();
await Promise.allSettled(workers.map((worker) => worker.close()));
await closeQueues();
await closeRedis();
await closeDatabase();
clearTimeout(forceTimer);
process.exit(exitCode);
};
process.once("SIGTERM", () => shutdown("SIGTERM"));
process.once("SIGINT", () => shutdown("SIGINT"));
process.once("unhandledRejection", (reason) => {
const error = reason instanceof Error ? reason : new Error("Unhandled rejection");
console.error("Unhandled worker rejection", { name: error.name, message: error.message });
shutdown("unhandledRejection", 1);
});
process.once("uncaughtException", (error) => {
console.error("Uncaught worker exception", { name: error.name, message: error.message });
shutdown("uncaughtException", 1);
});
};
if (require.main === module) {
startWorkers().catch((error) => {
console.error("Worker startup failed", { name: error.name, message: error.message });
process.exitCode = 1;
});
}
module.exports = { startWorkers };