feat: stabilize API startup and lifecycle management
- Refactor server initialization to separate concerns and improve error handling. - Implement centralized environment validation using Zod. - Introduce database, Redis, and queue lifecycle management. - Add health check endpoints for liveness and readiness. - Enhance error handling middleware for better response structure. - Implement rate limiting for API endpoints. - Add request ID middleware for traceability. - Create Sequelize CLI configuration and baseline migration for schema management. - Establish CI workflow with Gitea for testing and syntax checks. - Document foundational changes and migration strategy in PHASE_0_FOUNDATION_STABILIZATION.md. - Add Docker Compose configuration for local development and testing. - Implement unit and integration tests for critical functionality.
This commit is contained in:
@@ -14,16 +14,18 @@ const { ExpressAdapter } = require("@bull-board/express");
|
||||
const { BullMQAdapter } = require("@bull-board/api/bullMQAdapter");
|
||||
|
||||
const activityQueue = require("../queues/activity.queue");
|
||||
const documentQueue = require("../queues/document.queue");
|
||||
const logQueue = require("../queues/log.queue");
|
||||
|
||||
const serverAdapter = new ExpressAdapter();
|
||||
serverAdapter.setBasePath("/admin/queues");
|
||||
|
||||
const { addQueue, removeQueue, setQueues, replaceQueues } =
|
||||
createBullBoard({
|
||||
queues: [new BullMQAdapter(activityQueue)],
|
||||
queues: [activityQueue, documentQueue, logQueue].map((queue) => new BullMQAdapter(queue)),
|
||||
serverAdapter,
|
||||
});
|
||||
|
||||
module.exports = {
|
||||
bullBoardRouter: serverAdapter.getRouter(),
|
||||
};
|
||||
};
|
||||
|
||||
@@ -0,0 +1,9 @@
|
||||
const db = require("../models");
|
||||
|
||||
const initializeDatabase = async () => db.sequelize.authenticate();
|
||||
const checkDatabase = async () => {
|
||||
try { await db.sequelize.authenticate(); return true; } catch (_error) { return false; }
|
||||
};
|
||||
const closeDatabase = async () => db.sequelize.close();
|
||||
|
||||
module.exports = { initializeDatabase, checkDatabase, closeDatabase };
|
||||
@@ -12,10 +12,10 @@
|
||||
require("dotenv").config();
|
||||
|
||||
module.exports = {
|
||||
HOST: process.env.DB_HOST || "localhost",
|
||||
USER: process.env.DB_USER || "root",
|
||||
PASSWORD: process.env.DB_PASSWORD || "",
|
||||
DB: process.env.DB_NAME || "oceanic-db",
|
||||
HOST: process.env.DB_HOST,
|
||||
USER: process.env.DB_USER,
|
||||
PASSWORD: process.env.DB_PASSWORD,
|
||||
DB: process.env.DB_NAME,
|
||||
PORT: process.env.DB_PORT || 3306,
|
||||
DIALECT: "mysql",
|
||||
|
||||
|
||||
@@ -0,0 +1,64 @@
|
||||
const { z } = require("zod");
|
||||
|
||||
const booleanString = z.enum(["true", "false"]).default("false").transform((value) => value === "true");
|
||||
|
||||
const envSchema = z.object({
|
||||
NODE_ENV: z.enum(["development", "test", "production"]).default("development"),
|
||||
PORT: z.coerce.number().int().min(1).max(65535).default(3070),
|
||||
DB_HOST: z.string().min(1),
|
||||
DB_PORT: z.coerce.number().int().min(1).max(65535).default(3306),
|
||||
DB_NAME: z.string().min(1),
|
||||
DB_USER: z.string().min(1),
|
||||
DB_PASSWORD: z.string(),
|
||||
JWT_SECRET: z.string().min(32, "JWT_SECRET must contain at least 32 characters"),
|
||||
REFRESH_TOKEN_SECRET: z.string().min(32, "REFRESH_TOKEN_SECRET must contain at least 32 characters"),
|
||||
REDIS_HOST: z.string().min(1),
|
||||
REDIS_PORT: z.coerce.number().int().min(1).max(65535).default(6379),
|
||||
REDIS_PASSWORD: z.string().optional(),
|
||||
FRONTEND_URL: z.string().url(),
|
||||
TRUST_PROXY: z.coerce.number().int().min(0).max(10).default(0),
|
||||
JSON_BODY_LIMIT: z.string().default("1mb"),
|
||||
API_RATE_LIMIT_WINDOW_MS: z.coerce.number().int().positive().default(900000),
|
||||
API_RATE_LIMIT_MAX: z.coerce.number().int().positive().default(300),
|
||||
SENSITIVE_RATE_LIMIT_WINDOW_MS: z.coerce.number().int().positive().default(900000),
|
||||
SENSITIVE_RATE_LIMIT_MAX: z.coerce.number().int().positive().default(20),
|
||||
RUN_CRON: booleanString,
|
||||
ENABLE_MAIL: booleanString,
|
||||
ENABLE_S3: booleanString,
|
||||
SHUTDOWN_TIMEOUT_MS: z.coerce.number().int().positive().default(10000),
|
||||
CACHE: booleanString,
|
||||
MAIL_HOST: z.string().min(1).optional(), MAIL_PORT: z.coerce.number().int().positive().optional(),
|
||||
MAIL_SECURE: z.enum(["true", "false"]).optional(), MAIL_USER: z.string().optional(),
|
||||
MAIL_PASS: z.string().optional(), MAIL_FROM: z.string().optional(),
|
||||
AWS_REGION: z.string().optional(), AWS_ACCESS_KEY_ID: z.string().optional(),
|
||||
AWS_SECRET_ACCESS_KEY: z.string().optional(), AWS_S3_BUCKET_NAME: z.string().optional(),
|
||||
DOCS_USER: z.string().optional(), DOCS_PASS: z.string().optional(),
|
||||
}).superRefine((env, context) => {
|
||||
const requireFeature = (enabled, names) => {
|
||||
if (!enabled) return;
|
||||
for (const name of names) {
|
||||
if (!env[name]) context.addIssue({ code: "custom", path: [name], message: `${name} is required when enabled` });
|
||||
}
|
||||
};
|
||||
requireFeature(env.ENABLE_MAIL, ["MAIL_HOST", "MAIL_PORT", "MAIL_USER", "MAIL_PASS", "MAIL_FROM"]);
|
||||
requireFeature(env.ENABLE_S3, ["AWS_REGION", "AWS_ACCESS_KEY_ID", "AWS_SECRET_ACCESS_KEY", "AWS_S3_BUCKET_NAME"]);
|
||||
});
|
||||
|
||||
let validatedEnv;
|
||||
const validateEnvironment = (source = process.env) => {
|
||||
const result = envSchema.safeParse(source);
|
||||
if (!result.success) {
|
||||
const names = [...new Set(result.error.issues.map((issue) => issue.path.join(".") || "environment"))];
|
||||
throw new Error(`Invalid environment configuration: ${names.join(", ")}`);
|
||||
}
|
||||
validatedEnv = result.data;
|
||||
return validatedEnv;
|
||||
};
|
||||
const getEnvironment = () => validatedEnv || validateEnvironment();
|
||||
const getOptionalFeatureStatus = (env = process.env) => ({
|
||||
mail: Boolean(env.MAIL_HOST && env.MAIL_PORT && env.MAIL_USER && env.MAIL_PASS && env.MAIL_FROM),
|
||||
s3: Boolean(env.AWS_REGION && env.AWS_ACCESS_KEY_ID && env.AWS_SECRET_ACCESS_KEY && env.AWS_S3_BUCKET_NAME),
|
||||
docs: Boolean(env.DOCS_USER && env.DOCS_PASS),
|
||||
});
|
||||
|
||||
module.exports = { validateEnvironment, getEnvironment, getOptionalFeatureStatus };
|
||||
@@ -0,0 +1,9 @@
|
||||
const activityQueue = require("../queues/activity.queue");
|
||||
const documentQueue = require("../queues/document.queue");
|
||||
const logQueue = require("../queues/log.queue");
|
||||
|
||||
const closeQueues = async () => {
|
||||
await Promise.allSettled([activityQueue.close(), documentQueue.close(), logQueue.close()]);
|
||||
};
|
||||
|
||||
module.exports = { closeQueues };
|
||||
@@ -11,13 +11,14 @@
|
||||
|
||||
const { Redis } = require("ioredis");
|
||||
|
||||
const createRedisConnection = () => {
|
||||
const createRedisConnection = (options = {}) => {
|
||||
const redis = new Redis({
|
||||
host: process.env.REDIS_HOST || "redis",
|
||||
port: process.env.REDIS_PORT || 6379,
|
||||
password: process.env.REDIS_PASSWORD,
|
||||
maxRetriesPerRequest: null,
|
||||
enableReadyCheck: false,
|
||||
lazyConnect: options.lazyConnect ?? true,
|
||||
});
|
||||
|
||||
// Connection events
|
||||
|
||||
@@ -0,0 +1,14 @@
|
||||
const redis = require("./redisClient");
|
||||
|
||||
const initializeRedis = async () => {
|
||||
if (redis.status === "wait") await redis.connect();
|
||||
if (redis.status !== "ready") await redis.ping();
|
||||
};
|
||||
const checkRedis = async () => {
|
||||
try { return (await redis.ping()) === "PONG"; } catch (_error) { return false; }
|
||||
};
|
||||
const closeRedis = async () => {
|
||||
if (redis.status !== "end") await redis.quit();
|
||||
};
|
||||
|
||||
module.exports = { initializeRedis, checkRedis, closeRedis };
|
||||
@@ -12,6 +12,6 @@
|
||||
|
||||
const createRedisConnection = require("./redis.config");
|
||||
|
||||
const redis = createRedisConnection();
|
||||
const redis = createRedisConnection({ lazyConnect: true });
|
||||
|
||||
module.exports = redis;
|
||||
module.exports = redis;
|
||||
|
||||
+10
-21
@@ -1,22 +1,11 @@
|
||||
/**
|
||||
* 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.
|
||||
*/
|
||||
const { S3Client } = require("@aws-sdk/client-s3");
|
||||
|
||||
// app/config/s3.config.js
|
||||
|
||||
// const { S3Client } = require("@aws-sdk/client-s3");
|
||||
|
||||
// const s3 = new S3Client({
|
||||
// region: process.env.AWS_REGION,
|
||||
// credentials: {
|
||||
// accessKeyId: process.env.AWS_ACCESS_KEY_ID,
|
||||
// secretAccessKey: process.env.AWS_SECRET_ACCESS_KEY,
|
||||
// },
|
||||
// });
|
||||
|
||||
// module.exports = s3;
|
||||
module.exports = process.env.ENABLE_S3 === "true"
|
||||
? new S3Client({
|
||||
region: process.env.AWS_REGION,
|
||||
credentials: {
|
||||
accessKeyId: process.env.AWS_ACCESS_KEY_ID,
|
||||
secretAccessKey: process.env.AWS_SECRET_ACCESS_KEY,
|
||||
},
|
||||
})
|
||||
: { send: async () => { throw new Error("S3 functionality is not enabled"); } };
|
||||
|
||||
@@ -0,0 +1,17 @@
|
||||
require("dotenv").config();
|
||||
|
||||
const required = ["DB_HOST", "DB_NAME", "DB_USER"];
|
||||
const missing = required.filter((name) => !process.env[name]);
|
||||
if (missing.length) throw new Error(`Missing migration environment variables: ${missing.join(", ")}`);
|
||||
|
||||
const configuration = {
|
||||
username: process.env.DB_USER,
|
||||
password: process.env.DB_PASSWORD || "",
|
||||
database: process.env.DB_NAME,
|
||||
host: process.env.DB_HOST,
|
||||
port: Number(process.env.DB_PORT || 3306),
|
||||
dialect: "mysql",
|
||||
logging: false,
|
||||
};
|
||||
|
||||
module.exports = { development: configuration, test: configuration, production: configuration };
|
||||
@@ -0,0 +1,46 @@
|
||||
const { ValidationError, UniqueConstraintError } = require("sequelize");
|
||||
const { ZodError } = require("zod");
|
||||
|
||||
class AppError extends Error {
|
||||
constructor(status, code, message, details) {
|
||||
super(message);
|
||||
this.status = status;
|
||||
this.code = code;
|
||||
this.details = details;
|
||||
}
|
||||
}
|
||||
|
||||
const notFound = (req, _res, next) => next(new AppError(404, "NOT_FOUND", "Route not found"));
|
||||
|
||||
const errorHandler = (error, req, res, _next) => {
|
||||
let status = error.status || error.statusCode || 500;
|
||||
let code = error.code || "INTERNAL_ERROR";
|
||||
let message = error.message || "An unexpected error occurred";
|
||||
let details = error.details;
|
||||
|
||||
if (error instanceof ZodError) {
|
||||
status = 400; code = "VALIDATION_ERROR"; message = "Invalid request data";
|
||||
details = error.issues.map(({ path, message: detailMessage }) => ({ field: path.join("."), message: detailMessage }));
|
||||
} else if (error instanceof UniqueConstraintError) {
|
||||
status = 409; code = "CONFLICT"; message = "A record with these values already exists";
|
||||
details = error.errors?.map(({ path, message: detailMessage }) => ({ field: path, message: detailMessage }));
|
||||
} else if (error instanceof ValidationError) {
|
||||
status = 400; code = "VALIDATION_ERROR"; message = "Invalid request data";
|
||||
details = error.errors?.map(({ path, message: detailMessage }) => ({ field: path, message: detailMessage }));
|
||||
} else if (error.type === "entity.too.large") {
|
||||
status = 413; code = "PAYLOAD_TOO_LARGE"; message = "Request body is too large";
|
||||
} else if (error.name === "UnauthorizedError" || error.name === "JsonWebTokenError") {
|
||||
status = 401; code = "UNAUTHORIZED"; message = "Authentication failed";
|
||||
}
|
||||
|
||||
if (status >= 500) {
|
||||
console.error(`[${req.id || "no-request-id"}] Request failed`, { name: error.name, message: error.message });
|
||||
if (process.env.NODE_ENV === "production") message = "An unexpected error occurred";
|
||||
}
|
||||
|
||||
const payload = { success: false, error: { code, message }, requestId: req.id };
|
||||
if (details && status < 500) payload.error.details = details;
|
||||
res.status(status).json(payload);
|
||||
};
|
||||
|
||||
module.exports = { AppError, notFound, errorHandler };
|
||||
@@ -0,0 +1,21 @@
|
||||
const { rateLimit } = require("express-rate-limit");
|
||||
|
||||
const response = { success: false, error: { code: "RATE_LIMITED", message: "Too many requests" } };
|
||||
|
||||
const generalApiLimiter = rateLimit({
|
||||
windowMs: Number(process.env.API_RATE_LIMIT_WINDOW_MS) || 15 * 60 * 1000,
|
||||
limit: Number(process.env.API_RATE_LIMIT_MAX) || 300,
|
||||
standardHeaders: "draft-8",
|
||||
legacyHeaders: false,
|
||||
message: response,
|
||||
});
|
||||
|
||||
const sensitiveLimiter = rateLimit({
|
||||
windowMs: Number(process.env.SENSITIVE_RATE_LIMIT_WINDOW_MS) || 15 * 60 * 1000,
|
||||
limit: Number(process.env.SENSITIVE_RATE_LIMIT_MAX) || 20,
|
||||
standardHeaders: "draft-8",
|
||||
legacyHeaders: false,
|
||||
message: response,
|
||||
});
|
||||
|
||||
module.exports = { generalApiLimiter, sensitiveLimiter };
|
||||
@@ -0,0 +1,10 @@
|
||||
const { randomUUID } = require("crypto");
|
||||
|
||||
const SAFE_REQUEST_ID = /^[A-Za-z0-9_-]{8,128}$/;
|
||||
|
||||
module.exports = (req, res, next) => {
|
||||
const supplied = req.get("x-request-id");
|
||||
req.id = supplied && SAFE_REQUEST_ID.test(supplied) ? supplied : randomUUID();
|
||||
res.setHeader("X-Request-ID", req.id);
|
||||
next();
|
||||
};
|
||||
+2
-6
@@ -13,11 +13,7 @@
|
||||
const { Sequelize, DataTypes } = require("sequelize");
|
||||
const dbConfig = require("../config/db.config");
|
||||
|
||||
let loggingOption = false;
|
||||
|
||||
if(process.env.NODE_ENV === 'production') {
|
||||
loggingOption = true;
|
||||
}
|
||||
const loggingOption = process.env.NODE_ENV === "development" ? console.log : false;
|
||||
|
||||
const sequelize = new Sequelize(
|
||||
dbConfig.DB,
|
||||
@@ -76,4 +72,4 @@ Object.keys(db).forEach(model => {
|
||||
|
||||
|
||||
|
||||
module.exports = db;
|
||||
module.exports = db;
|
||||
|
||||
@@ -13,6 +13,9 @@ const express = require('express');
|
||||
const router = express.Router();
|
||||
const authController = require('../controllers/auth.controller');
|
||||
const {authenticate} = require('../middleware/auth.middleware');
|
||||
const { sensitiveLimiter } = require('../middleware/rateLimit.middleware');
|
||||
|
||||
router.use(sensitiveLimiter);
|
||||
|
||||
// GET /api/auth/me
|
||||
router.get("/me", authenticate, (req, res) => {
|
||||
|
||||
@@ -6,6 +6,7 @@ const express = require("express");
|
||||
const fs = require("fs");
|
||||
const path = require("path");
|
||||
const docsSession = require("../middleware/docsSession.middleware");
|
||||
const { sensitiveLimiter } = require("../middleware/rateLimit.middleware");
|
||||
|
||||
const router = express.Router();
|
||||
|
||||
@@ -47,7 +48,7 @@ router.get("/view/:file", docsSession, (req, res) => {
|
||||
});
|
||||
|
||||
// Docs login (simple)
|
||||
router.post("/login", (req, res) => {
|
||||
router.post("/login", sensitiveLimiter, (req, res) => {
|
||||
const { username, password } = req.body;
|
||||
|
||||
if (
|
||||
|
||||
@@ -0,0 +1,24 @@
|
||||
const express = require("express");
|
||||
const { checkDatabase } = require("../config/database.lifecycle");
|
||||
const { checkRedis } = require("../config/redis.lifecycle");
|
||||
|
||||
const router = express.Router();
|
||||
|
||||
const live = (_req, res) => res.json({
|
||||
status: "ok",
|
||||
service: "zumri-api",
|
||||
timestamp: new Date().toISOString(),
|
||||
uptime: process.uptime(),
|
||||
});
|
||||
|
||||
router.get("/live", live);
|
||||
router.get("/ready", async (_req, res) => {
|
||||
const [database, redis] = await Promise.all([checkDatabase(), checkRedis()]);
|
||||
const ready = database && redis;
|
||||
res.status(ready ? 200 : 503).json({
|
||||
status: ready ? "ready" : "not_ready",
|
||||
checks: { database: database ? "ok" : "error", redis: redis ? "ok" : "error" },
|
||||
});
|
||||
});
|
||||
|
||||
module.exports = { healthRouter: router, live };
|
||||
@@ -19,10 +19,11 @@ const {
|
||||
checkPermission,
|
||||
} = require("../middleware/permission.middleware");
|
||||
const PERMISSIONS = require("../constants/permissions");
|
||||
const { sensitiveLimiter } = require("../middleware/rateLimit.middleware");
|
||||
|
||||
router.post("/req-reset-password", profileController.requestPasswordReset);
|
||||
router.post("/req-reset-password", sensitiveLimiter, profileController.requestPasswordReset);
|
||||
|
||||
router.post("/reset-password", profileController.resetPassword);
|
||||
router.post("/reset-password", sensitiveLimiter, profileController.resetPassword);
|
||||
|
||||
router.post(
|
||||
"/change-password",
|
||||
|
||||
+10
-6
@@ -12,19 +12,23 @@
|
||||
const jwt = require("jsonwebtoken");
|
||||
require("dotenv").config();
|
||||
|
||||
const JWT_SECRET = process.env.JWT_SECRET || "your_jwt_secret_key";
|
||||
const JWT_EXPIRES_IN = process.env.JWT_EXPIRES_IN || "15m"; // token validity
|
||||
|
||||
const REFRESH_TOKEN_SECRET = process.env.REFRESH_TOKEN_SECRET || "your_refresh_token_secret_key";
|
||||
const REFRESH_TOKEN_DAYS = process.env.REFRESH_TOKEN_DAYS || "7d"; // refresh token validity
|
||||
|
||||
const getSecret = (name) => {
|
||||
const value = process.env[name];
|
||||
if (!value || value.length < 32) throw new Error(`${name} is not configured securely`);
|
||||
return value;
|
||||
};
|
||||
|
||||
/**
|
||||
* Generate JWT token
|
||||
* @param {Object} payload - usually { id, email, role }
|
||||
* @returns string
|
||||
*/
|
||||
const generateToken = (payload) => {
|
||||
return jwt.sign(payload, JWT_SECRET, { expiresIn: JWT_EXPIRES_IN });
|
||||
return jwt.sign(payload, getSecret("JWT_SECRET"), { expiresIn: JWT_EXPIRES_IN });
|
||||
};
|
||||
|
||||
/**
|
||||
@@ -33,15 +37,15 @@ const generateToken = (payload) => {
|
||||
* @returns payload or throws error
|
||||
*/
|
||||
const verifyToken = (token) => {
|
||||
return jwt.verify(token, JWT_SECRET);
|
||||
return jwt.verify(token, getSecret("JWT_SECRET"));
|
||||
};
|
||||
|
||||
const generateRefreshToken = (payload) => {
|
||||
return jwt.sign(payload, REFRESH_TOKEN_SECRET, { expiresIn: REFRESH_TOKEN_DAYS });
|
||||
return jwt.sign(payload, getSecret("REFRESH_TOKEN_SECRET"), { expiresIn: REFRESH_TOKEN_DAYS });
|
||||
}
|
||||
|
||||
const verifyRefreshToken = (token) => {
|
||||
return jwt.verify(token, REFRESH_TOKEN_SECRET);
|
||||
return jwt.verify(token, getSecret("REFRESH_TOKEN_SECRET"));
|
||||
}
|
||||
|
||||
module.exports = { generateToken, verifyToken, generateRefreshToken, verifyRefreshToken };
|
||||
|
||||
@@ -23,10 +23,12 @@ const transporter = nodemailer.createTransport({
|
||||
auth: mailConfig.auth,
|
||||
});
|
||||
|
||||
transporter.verify((err) => {
|
||||
if (err) console.error("Mail server connection failed", err);
|
||||
else console.log("Mail server ready");
|
||||
});
|
||||
if (process.env.NODE_ENV !== "test" && process.env.ENABLE_MAIL === "true") {
|
||||
transporter.verify((err) => {
|
||||
if (err) console.error("Mail server connection failed", { message: err.message });
|
||||
else console.log("Mail server ready");
|
||||
});
|
||||
}
|
||||
|
||||
// Utility to load template and replace placeholders
|
||||
const loadTemplate = (templateName, variables = {}) => {
|
||||
@@ -42,6 +44,9 @@ const loadTemplate = (templateName, variables = {}) => {
|
||||
};
|
||||
|
||||
const send = async ({ to, subject, templateName, templateVars = {}, text }) => {
|
||||
if (process.env.ENABLE_MAIL !== "true") {
|
||||
throw new Error("Email functionality is not enabled");
|
||||
}
|
||||
const html = templateName ? loadTemplate(templateName, templateVars) : undefined;
|
||||
|
||||
const mailOptions = {
|
||||
|
||||
@@ -15,8 +15,7 @@ const { getSignedUrl } = require("@aws-sdk/s3-request-presigner");
|
||||
const s3 = require("../config/s3.config");
|
||||
const { v4: uuidv4 } = require("uuid");
|
||||
const path = require("path");
|
||||
const createRedisConnection = require("../config/redis.config");
|
||||
const redis = createRedisConnection();
|
||||
const redis = require("../config/redisClient");
|
||||
const { log } = require("./consoleLog.utill");
|
||||
|
||||
// Upload + return key (BEST PRACTICE)
|
||||
|
||||
+55
-26
@@ -1,34 +1,63 @@
|
||||
/**
|
||||
* 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 { validateEnvironment } = require("../config/env.config");
|
||||
|
||||
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 () => {
|
||||
const startWorkers = async () => {
|
||||
const env = validateEnvironment();
|
||||
const { initializeDatabase, closeDatabase } = require("../config/database.lifecycle");
|
||||
const { initializeRedis, closeRedis } = require("../config/redis.lifecycle");
|
||||
try {
|
||||
createActivityWorker();
|
||||
createLogWorker();
|
||||
await createDocumentWorker();
|
||||
|
||||
console.log("✅ All workers started");
|
||||
} catch (err) {
|
||||
console.error("❌ Error starting workers:", err.message);
|
||||
process.exit(1);
|
||||
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 workers = [createActivityWorker(), createLogWorker(), await createDocumentWorker()];
|
||||
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 };
|
||||
|
||||
Reference in New Issue
Block a user