From b6b345f245dea8424b2b35d9bb62ef1bc3be319c Mon Sep 17 00:00:00 2001 From: Sathira Sri Sathara Date: Thu, 3 Sep 2026 14:22:34 +0530 Subject: [PATCH] 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. --- .env.sample | 10 + Documentation/API_CROSS_CUTTING_SERVICES.md | 31 ++ Documentation/CURRENT_BACKEND_STATUS.md | 8 + .../PHASE_2_CROSS_CUTTING_SERVICES.md | 93 ++++ app/config/bullBoard.config.js | 3 +- app/config/env.config.js | 7 + app/config/queue.lifecycle.js | 3 +- app/config/s3.config.js | 2 + app/constants/permissions.js | 6 +- .../generateDocument.controller.js | 432 +----------------- app/controllers/notification.controller.js | 150 +----- app/controllers/profile.controller.js | 10 +- app/controllers/upload.controller.js | 159 ++----- app/logic/documents/registry.js | 6 +- app/middleware/docsSession.middleware.js | 27 +- app/middleware/upload.middleware.js | 19 +- app/models/activities/userActivities.model.js | 9 +- app/models/document/document.model.js | 12 +- app/models/document/documentType.model.js | 4 +- app/models/index.js | 1 + .../notificationDelivery.model.js | 11 + .../notification/userNotification.model.js | 2 + app/models/upload/upload.model.js | 9 + app/queues/activity.queue.js | 3 +- app/queues/document.queue.js | 2 +- app/queues/email.queue.js | 3 + app/queues/log.queue.js | 4 +- app/routes/activity.routes.js | 4 +- app/routes/docs.routes.js | 27 +- app/routes/document.routes.js | 106 +---- app/routes/notification.routes.js | 80 +--- app/routes/profile.routes.js | 10 + app/routes/upload.routes.js | 5 +- app/services/activity.service.js | 24 +- app/services/auth/email.service.js | 6 +- app/services/email/email.service.js | 15 + .../notification-policy.service.js | 7 + .../storage/file-validation.service.js | 19 + app/services/storage/storage.service.js | 19 + app/utils/consoleLog.utill.js | 6 +- app/utils/documentJob.util.js | 2 +- app/utils/emailVerification.util.js | 17 +- app/utils/mail.util.js | 6 +- app/utils/passwordReset.utill.js | 36 +- app/utils/s3Upload.utill.js | 88 +--- app/workers/activity.worker.js | 3 +- app/workers/document.worker.js | 114 +---- app/workers/email.worker.js | 17 + app/workers/index.js | 3 +- app/workers/log.worker.js | 6 +- cron/notificationCleaning.cron.js | 111 +---- ...03020000-phase-2-cross-cutting-services.js | 43 ++ tests/unit/cross-cutting-services.test.js | 20 + tests/unit/storage.service.test.js | 21 + 54 files changed, 593 insertions(+), 1248 deletions(-) create mode 100644 Documentation/API_CROSS_CUTTING_SERVICES.md create mode 100644 Documentation/PHASE_2_CROSS_CUTTING_SERVICES.md create mode 100644 app/models/notification/notificationDelivery.model.js create mode 100644 app/queues/email.queue.js create mode 100644 app/services/email/email.service.js create mode 100644 app/services/notification/notification-policy.service.js create mode 100644 app/services/storage/file-validation.service.js create mode 100644 app/services/storage/storage.service.js create mode 100644 app/workers/email.worker.js create mode 100644 migrations/20260903020000-phase-2-cross-cutting-services.js create mode 100644 tests/unit/cross-cutting-services.test.js create mode 100644 tests/unit/storage.service.test.js diff --git a/.env.sample b/.env.sample index 685dfc6..0a38a8b 100644 --- a/.env.sample +++ b/.env.sample @@ -52,6 +52,16 @@ AWS_ACCESS_KEY_ID= AWS_SECRET_ACCESS_KEY= AWS_REGION= AWS_S3_BUCKET_NAME= +S3_ENDPOINT= +S3_FORCE_PATH_STYLE=false +S3_SIGNED_URL_TTL_SECONDS=900 +S3_MAX_UPLOAD_BYTES=5242880 + +# Cross-cutting worker and retention settings +EMAIL_QUEUE_CONCURRENCY=5 +DOCUMENT_QUEUE_CONCURRENCY=2 +NOTIFICATION_RETENTION_DAYS=90 +LOG_RETENTION_DAYS=30 # Optional documentation login DOCS_USER= diff --git a/Documentation/API_CROSS_CUTTING_SERVICES.md b/Documentation/API_CROSS_CUTTING_SERVICES.md new file mode 100644 index 0000000..b5d3d43 --- /dev/null +++ b/Documentation/API_CROSS_CUTTING_SERVICES.md @@ -0,0 +1,31 @@ +# ZUMRI Cross-Cutting Services API + +All paths also exist below `/api`; clients should use `/api/v1`. Protected routes accept the Phase 1 access cookie or bearer token. Examples use placeholders and never expose storage keys. + +## Media + +- `POST /api/v1/upload` — multipart field `file`, optional text field `use_for`; creates a private owner-bound upload. +- `GET /api/v1/upload/signed-url/:id` — returns `{ id, url, expiresIn }` after ownership/permission checks. +- `DELETE /api/v1/upload/:id` — marks an owned/authorized upload deleted and removes its object best-effort. +- `GET /api/v1/profile/me/avatar` and `/background` — signed self profile media access. Legacy owner-checked ID routes remain. + +## Notifications + +- `POST /api/v1/notification` — `notifications.manage`; body includes headline, description, `USER|ANNOUNCEMENT`, and `userIds` for USER messages. +- `GET /api/v1/notification/announcements` +- `GET /api/v1/notification/me` +- `PATCH /api/v1/notification/:notificationId/read` +- `PATCH /api/v1/notification/read-all` + +## Documents + +- `GET /api/v1/document/types` +- `GET /api/v1/document/saved` +- `POST /api/v1/document/draft` +- `POST /api/v1/document/generate` with `{ "document":"", "documentType":"pdf", "documentData":{} }`; returns HTTP 202 and persisted status. +- `GET /api/v1/document/jobs/:jobId` +- `GET /api/v1/document/:docId` +- `GET /api/v1/document/:docId/download` — returns a short-lived URL; it does not delete the artifact. +- `DELETE /api/v1/document/job/:jobId` + +The legacy reference-number GET returns 410 because reads must not consume sequences. References are assigned as part of resource creation. diff --git a/Documentation/CURRENT_BACKEND_STATUS.md b/Documentation/CURRENT_BACKEND_STATUS.md index cb6d00c..f3bc3b4 100644 --- a/Documentation/CURRENT_BACKEND_STATUS.md +++ b/Documentation/CURRENT_BACKEND_STATUS.md @@ -1,5 +1,13 @@ # ZUMRI Current Backend Status +## Phase 2 Completion Update + +Completion date: 2026-09-03. Phase 2 hardens the existing shared-service foundation without adding commerce domains. Module 14 (notifications) is now approximately 72%; Module 19 (file/media) 82%; Module 20 (audit/config/logging) 68%; and Module 21 (background jobs) 78%. The document subsystem is approximately 82%. + +Storage now has a reusable S3/S3-compatible boundary, actual-content validation, controlled keys, checksums, owner/status metadata, compensation, and authorization-safe signed URLs. Email is routed through BullMQ with delivery status and final-failure persistence. Notifications have fixed aliases, unique assignment migration, self-only inbox/read operations, announcement support, and preference policy. Queue defaults, idempotent job IDs, Bull Board registration, safe failure handling, stronger append-only activity records, and structured redacted logs are in place. Documents now have ownership, controlled registry validation, persisted generation lifecycle, safe status, and non-destructive signed download. + +The full mocked suite contains 10 suites/40 tests and passes; syntax checks cover 153 JavaScript files. The new migration was not executed. Remaining work is staging migration/data pre-checks plus real MySQL, Redis, S3-compatible, SMTP, PDF/browser, retention-volume, and concurrency validation. After those operational checks and permission seeding, it is safe to begin Phase 3. See `Documentation/PHASE_2_CROSS_CUTTING_SERVICES.md` and `Documentation/API_CROSS_CUTTING_SERVICES.md`. + Audit date: 2026-09-03 Scope: repository source, configuration, lockfile, existing documentation, safe syntax/test/dependency checks. No database, Redis, S3, email, or other external service was mutated. diff --git a/Documentation/PHASE_2_CROSS_CUTTING_SERVICES.md b/Documentation/PHASE_2_CROSS_CUTTING_SERVICES.md new file mode 100644 index 0000000..2f13a5b --- /dev/null +++ b/Documentation/PHASE_2_CROSS_CUTTING_SERVICES.md @@ -0,0 +1,93 @@ +# ZUMRI Phase 2 Cross-Cutting Services + +## Objective + +Harden the shared storage, messaging, queue, audit, logging, and document infrastructure without starting commerce modules. + +## Existing Components Reused + +The existing AWS SDK client, Upload/Notification/UserNotification/Document models, Redis/BullMQ topology, generators, templates, reference utility, workers, cron lifecycle, Phase 0 operations, and Phase 1 identity/RBAC remain the foundation. + +## Storage Architecture + +`storage.service.js` is the sole AWS SDK boundary. It supports AWS S3 and endpoint/path-style compatible providers, buffer upload, delete, HEAD existence checks, and signed downloads. Objects remain private. + +## Upload Security + +Multer performs an early allowlist/size check. The service then rejects empty content and validates JPEG, PNG, WebP, PDF, and XLSX magic bytes against the claimed MIME. It derives the extension, sanitizes display filenames, and stores SHA-256 checksums. Object keys contain a controlled owner identifier, UTC year/month, UUID, and detected extension; original filenames and personal data are excluded. + +## File Ownership + +Uploads record uploader, owner type/id, purpose, visibility, and lifecycle status. Signed URL and deletion endpoints load metadata and enforce owner or explicit permission access. DB persistence failure after upload triggers best-effort object deletion. Deletion marks metadata DELETED before object removal. + +## Signed URL Policy + +Downloads use `S3_SIGNED_URL_TTL_SECONDS` (default 900 seconds). URLs are generated after each authorization decision and are not globally cached or exposed with bucket/key details. + +## Email Architecture + +Auth email enters `email.service.js`, creates a minimal delivery record, and queues a template-keyed job. The worker alone calls Nodemailer. OTPs/links may exist transiently in job data, so jobs have aggressive completion retention and payloads must never be logged. + +## Email Retry Strategy + +Email uses five exponential attempts. Envelope/message and SMTP 5xx failures are treated as permanent; transient provider/network errors retry. Final state is persisted without storing message bodies or variables. + +## Notification Architecture + +Notification aliases are explicit. Admin publication can assign one or many users atomically; announcements need no join rows. Self-service listing/read/read-all always derives the user from the access token. + +## Notification Preferences + +The central policy permits mandatory login OTP, password reset/change, and email verification even when marketing notifications are disabled. Optional external-channel messages honor profile preferences. In-app messages remain available. + +## Queue Architecture + +Activity, log, document, and email queues define retry, backoff, success retention, and failure retention appropriate to each workload. Bull Board includes all four and retains Phase 0 admin protection. + +## Idempotency + +Activity uses an event ID, email uses event/delivery ID, and documents use the generation record ID as BullMQ job ID. Workers check persisted state where duplicate execution could create a second artifact. + +## Failed Job Handling + +Email and document final failures update their associated database record with a bounded error code and timestamp. Stack traces and job payloads are not returned by APIs. + +## Audit Logging + +Activity events now support event ID, nullable actor, target, action/type, request/IP/user-agent context, sanitized JSON metadata, and occurrence time. APIs expose read operations only; inserts are idempotent by event ID. + +## Logging Security + +Queued file logs are JSON lines. Error stacks, request bodies, Authorization/Cookie values, and fields named like passwords, OTPs, tokens, or secrets are excluded/redacted. Log jobs have bounded retention; `LOG_RETENTION_DAYS` documents the intended operational file-retention window. + +## Document Generation Lifecycle + +Generation validates a controlled registry and PDF/XLSX format before queueing, creates an owner-bound record, and moves through QUEUED, PROCESSING, COMPLETED, or FAILED. Successful output becomes an owned Upload. Downloads create a signed URL and never delete the object. + +## Document Ownership + +Creator and owner may view status/data/download. SUPER_ADMIN bypasses; other administrative access requires the appropriate document permission. Cancellation follows the same ownership boundary. + +## Migrations + +`20260903020000-phase-2-cross-cutting-services.js` is new and forward-only. Before execution, back up and test a restored database. Pre-check duplicate `(user_id, notification_id)`, duplicate activity event IDs, duplicate document job IDs, duplicate document type names, and legacy uploads/documents without resolvable owners. Resolve duplicates explicitly; the migration intentionally does not delete data. + +## Permissions + +Shared names are `media.read/upload/delete`, `documents.read/create/delete`, `notifications.manage/read`, `audit.read`, and `queues.read`. SUPER_ADMIN retains the Phase 1 bypass. Production permission rows/grants must be seeded through the environment's controlled authorization process. + +## Environment Variables + +Added: `S3_ENDPOINT`, `S3_FORCE_PATH_STYLE`, `S3_SIGNED_URL_TTL_SECONDS`, `S3_MAX_UPLOAD_BYTES`, `EMAIL_QUEUE_CONCURRENCY`, `DOCUMENT_QUEUE_CONCURRENCY`, `NOTIFICATION_RETENTION_DAYS`, and `LOG_RETENTION_DAYS`. Optional integrations remain optional unless enabled. + +## Tests + +Unit coverage verifies magic bytes, mismatch/empty rejection, checksums, safe keys, HTML escaping, security-notification preference policy, and email queue retry/retention. Existing Phase 0/1 tests remain in the full suite. AWS, SMTP, Redis, and MySQL are not contacted. + +## Remaining Known Issues + +The migration has not been run against staging data. Real S3-compatible provider, SMTP, MySQL migration, Redis concurrency, large notification retention, and actual PDF/browser generation need staging validation. File-log deletion/rotation still belongs to deployment logrotate or a future controlled maintenance worker. Push delivery is intentionally an architecture placeholder only. + +## Phase 3 Prerequisites + +Complete the documented data pre-checks, apply all migrations to a restored database, seed shared permissions, run API and worker processes against staging Redis/MySQL, and exercise one upload/email/document lifecycle with non-production provider credentials. diff --git a/app/config/bullBoard.config.js b/app/config/bullBoard.config.js index fbcec22..6f2e0cb 100644 --- a/app/config/bullBoard.config.js +++ b/app/config/bullBoard.config.js @@ -16,13 +16,14 @@ 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 emailQueue = require("../queues/email.queue"); const serverAdapter = new ExpressAdapter(); serverAdapter.setBasePath("/admin/queues"); const { addQueue, removeQueue, setQueues, replaceQueues } = createBullBoard({ - queues: [activityQueue, documentQueue, logQueue].map((queue) => new BullMQAdapter(queue)), + queues: [activityQueue, documentQueue, logQueue, emailQueue].map((queue) => new BullMQAdapter(queue)), serverAdapter, }); diff --git a/app/config/env.config.js b/app/config/env.config.js index 8e708d4..369c22c 100644 --- a/app/config/env.config.js +++ b/app/config/env.config.js @@ -38,6 +38,13 @@ const envSchema = z.object({ 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(), + S3_ENDPOINT: z.string().url().optional(), S3_FORCE_PATH_STYLE: booleanString, + S3_SIGNED_URL_TTL_SECONDS: z.coerce.number().int().min(60).max(86400).default(900), + S3_MAX_UPLOAD_BYTES: z.coerce.number().int().positive().default(5242880), + EMAIL_QUEUE_CONCURRENCY: z.coerce.number().int().positive().default(5), + DOCUMENT_QUEUE_CONCURRENCY: z.coerce.number().int().positive().default(2), + NOTIFICATION_RETENTION_DAYS: z.coerce.number().int().positive().default(90), + LOG_RETENTION_DAYS: z.coerce.number().int().positive().default(30), DOCS_USER: z.string().optional(), DOCS_PASS: z.string().optional(), GOOGLE_CLIENT_ID: z.string().optional(), APPLE_CLIENT_ID: z.string().optional(), }).superRefine((env, context) => { diff --git a/app/config/queue.lifecycle.js b/app/config/queue.lifecycle.js index 1f22616..e9aa1b5 100644 --- a/app/config/queue.lifecycle.js +++ b/app/config/queue.lifecycle.js @@ -1,9 +1,10 @@ const activityQueue = require("../queues/activity.queue"); const documentQueue = require("../queues/document.queue"); const logQueue = require("../queues/log.queue"); +const emailQueue = require("../queues/email.queue"); const closeQueues = async () => { - await Promise.allSettled([activityQueue.close(), documentQueue.close(), logQueue.close()]); + await Promise.allSettled([activityQueue.close(), documentQueue.close(), logQueue.close(), emailQueue.close()]); }; module.exports = { closeQueues }; diff --git a/app/config/s3.config.js b/app/config/s3.config.js index 7ec7341..ce81fa5 100644 --- a/app/config/s3.config.js +++ b/app/config/s3.config.js @@ -7,5 +7,7 @@ module.exports = process.env.ENABLE_S3 === "true" accessKeyId: process.env.AWS_ACCESS_KEY_ID, secretAccessKey: process.env.AWS_SECRET_ACCESS_KEY, }, + ...(process.env.S3_ENDPOINT ? { endpoint: process.env.S3_ENDPOINT } : {}), + forcePathStyle: process.env.S3_FORCE_PATH_STYLE === "true", }) : { send: async () => { throw new Error("S3 functionality is not enabled"); } }; diff --git a/app/constants/permissions.js b/app/constants/permissions.js index 9879db7..e7ca20d 100644 --- a/app/constants/permissions.js +++ b/app/constants/permissions.js @@ -13,4 +13,8 @@ module.exports = { FINANCE_BASIS: "finance.basis", PRECOST: "precost.precost", -}; \ No newline at end of file + MEDIA_READ: "media.read", MEDIA_UPLOAD: "media.upload", MEDIA_DELETE: "media.delete", + DOCUMENTS_READ: "documents.read", DOCUMENTS_CREATE: "documents.create", DOCUMENTS_DELETE: "documents.delete", + NOTIFICATIONS_MANAGE: "notifications.manage", NOTIFICATIONS_READ: "notifications.read", + AUDIT_READ: "audit.read", QUEUES_READ: "queues.read", +}; diff --git a/app/controllers/generateDocument.controller.js b/app/controllers/generateDocument.controller.js index 4f69daf..ae0301b 100644 --- a/app/controllers/generateDocument.controller.js +++ b/app/controllers/generateDocument.controller.js @@ -1,417 +1,21 @@ -/** - * 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/controllers/generateDocument.controller.js - -const fs = require("fs"); -const path = require("path"); -const documentQueue = require("../queues/document.queue"); -const { createDocumentData } = require("../utils/document.utill"); -const { - generateId, - generateDocumentReferenceNo, -} = require("../utils/idGen.util"); - -const { GetObjectCommand, DeleteObjectCommand } = require("@aws-sdk/client-s3"); - -const {log} = require("../utils/consoleLog.utill"); - -const s3 = require("../config/s3.config"); +const crypto = require("crypto"); +const { z } = require("zod"); const db = require("../models"); -const Document = db.Document; -const DocumentType = db.DocumentType; +const documentQueue = require("../queues/document.queue"); +const registry = require("../logic/documents/registry"); +const storage = require("../services/storage/storage.service"); -const normalizeDocumentData = (value) => { - if (!value) { - return {}; - } +const requestSchema = z.object({ document: z.string().min(1).max(64), documentType: z.enum(["pdf", "excel"]), documentData: z.record(z.string(), z.unknown()).optional(), data: z.record(z.string(), z.unknown()).optional() }).passthrough(); +const normalize = (value) => String(value).toLowerCase().replace(/[\s_-]+/g, ""); +const allowed = (user, doc, permission = "documents.read") => doc.owner_id === user.id || doc.created_by === user.id || user.accountType === "super_admin" || (user.permissions || []).includes(permission); +const publicDoc = (doc) => ({ id: doc.doc_id, referenceNo: doc.reference_no, documentType: doc.doc_type, status: doc.status, jobId: doc.job_id, generatedAt: doc.generated_at, failedAt: doc.failed_at, failureCode: doc.failure_code, createdAt: doc.createdAt }); - if (typeof value === "string") { - try { - return JSON.parse(value); - } catch (error) { - return {}; - } - } - - return value; -}; - -/** - * Get available document types - */ -exports.getAvailableDocumentTypes = async (req, res) => { - try { - - const doc_types = await DocumentType.findAll({ - attributes: ["doc_type_name"], - }); - - const types = doc_types.map((t) => t.doc_type_name); - - return res.json({ - success: true, - availableDocumentTypes: types, - note: "Use these document type names in your requests (case-insensitive)", - }); - - } catch (error) { - log("Error fetching document types:", error.message); - return res.status(500).json({ - success: false, - message: error.message, - }); - } -}; - -// Get Saved documents -exports.getSavedDocuments = async (req, res) => { - try { - const { documentType } = req.body; - - if (!documentType) { - return res.status(400).json({ - success: false, - message: "documentType query parameter is required", - }); - } - - // Fetch saved documents based on documentType - const savedDocuments = await Document.findAll({ - where: { doc_type: documentType.toUpperCase() }, - order: [["createdAt", "DESC"]], - exclude: ["data"], - }); - - // extract only necessary fields to return - const formattedDocuments = savedDocuments.map((doc) => ({ - doc_id: doc.doc_id, - reference_no: doc.reference_no, - doc_type: doc.doc_type, - status: doc.status, - createdAt: doc.createdAt, - updatedAt: doc.updatedAt, - })); - - return res.json({ - success: true, - savedDocuments: formattedDocuments, - }); - } catch (error) { - log("Error fetching saved documents:", error.message); - return res.status(500).json({ - success: false, - message: error.message, - }); - } -}; - -// Get specific document data by doc_id -exports.getDocumentData = async (req, res) => { - try { - const { docId } = req.params; - - if (!docId) { - return res.status(400).json({ - success: false, - message: "docId parameter is required", - }); - } - - const document = await Document.findOne({ - where: { doc_id: docId }, - }); - - if (!document) { - return res.status(404).json({ - success: false, - message: "Document not found", - }); - } - - return res.json({ - success: true, - document: { - doc_id: document.doc_id, - reference_no: document.reference_no, - doc_type: document.doc_type, - data: document.data, - status: document.status, - createdAt: document.createdAt, - updatedAt: document.updatedAt, - }, - }); - - } catch (error) { - log("Error fetching document data:", error.message); - return res.status(500).json({ - success: false, - message: error.message, - }); - } -} - -exports.generateReferenceNo = async (req, res) => { - try { - const { documentType } = req.params; - - if (!documentType) { - return res.status(400).json({ - success: false, - message: "documentType is required", - }); - } - - // Generate reference number - const reference_no = await generateDocumentReferenceNo(documentType); - - return res.json({ - success: true, - reference_no, - }); - } catch (error) { - log("Error generating reference number:", error.message); - return res.status(500).json({ - success: false, - message: error.message, - }); - } -}; - -exports.generateDraftDocument = async (req, res) => { - try { - const { document } = req.body; - const documentData = normalizeDocumentData(req.body.documentData || req.body.data || req.body); - - // Validation - if (!document || !documentData) { - return res.status(400).json({ - success: false, - message: "document and documentData are required", - }); - } - - let documentDetails; - - if (document !== "PRECOST") { - // Save document details before generating - documentDetails = await createDocumentData( - document, - documentData, - "DRAFT", - ); - } else { - return res.status(400).json({ - success: false, - message: - "Invalid document type, This document type is not allowed to be generated as draft", - }); - } - - return res.status(201).json({ - success: true, - message: "Draft document created successfully", - documentDetails, - }); - } catch (error) { - log("Error generating draft document:", error.message); - return res.status(500).json({ - success: false, - message: error.message, - }); - } -}; - -/** - * Generate document asynchronously - * Returns jobId immediately - */ -exports.generateDocument = async (req, res) => { - try { - const { document, documentType } = req.body; - const documentData = normalizeDocumentData(req.body.documentData || req.body.data || req.body); - - // Validation - if (!document || !documentType || !documentData) { - return res.status(400).json({ - success: false, - message: "document, documentType, and documentData are required", - }); - } - - console.log( - `📨 generateDocument request: document="${document}", documentType="${documentType}"`, - ); - - if (document !== "PRECOST") { - // Save document details before generating - // If doc_id exists, finalize (update existing); otherwise create as DRAFT - const status = documentData.doc_id ? "FINAL" : "DRAFT"; - await createDocumentData(document, documentData, status); - } - - // Add job to queue - const job = await documentQueue.add("generate-document", { - document, - documentType, - data: documentData, - }); - - log( - `📋 Document generation job queued: ${job.id} (document="${document}", type="${documentType}")`, - ); - - return res.status(202).json({ - success: true, - message: "Document generation started", - jobId: job.id, - }); - } catch (error) { - log("Error queuing document generation:", error.message); - return res.status(500).json({ - success: false, - message: error.message, - }); - } -}; - -/** - * Get job status and result - */ -exports.getJobStatus = async (req, res) => { - try { - const { jobId } = req.params; - - if (!jobId) { - return res.status(400).json({ - success: false, - message: "jobId is required", - }); - } - - // Get job from queue - const job = await documentQueue.getJob(jobId); - - if (!job) { - return res.status(404).json({ - success: false, - message: "Job not found", - }); - } - - // Get job state - const state = await job.getState(); - const result = job.returnvalue; - const failedReason = job.failedReason; - - return res.json({ - success: true, - jobId: job.id, - state, // "waiting" | "active" | "completed" | "failed" | "delayed" - result: state === "completed" ? result : null, - error: state === "failed" ? failedReason : null, - attempts: job.attemptsMade, - stacktrace: job.stacktrace, - }); - } catch (error) { - log("Error fetching job status:", error.message); - return res.status(500).json({ - success: false, - message: error.message, - }); - } -}; - -/** - * Download generated document - */ -exports.downloadDocument = async (req, res) => { - try { - const { uuid } = req.params; - - const key = `uploads/${uuid}.pdf`; - - const command = new GetObjectCommand({ - Bucket: process.env.AWS_S3_BUCKET_NAME, - Key: key, - }); - - const response = await s3.send(command); - - res.setHeader( - "Content-Type", - response.ContentType || "application/octet-stream", - ); - - res.setHeader("Content-Disposition", `attachment; filename="${uuid}.pdf"`); - - response.Body.pipe(res); - - res.on("finish", async () => { - try { - await s3.send( - new DeleteObjectCommand({ - Bucket: process.env.AWS_S3_BUCKET_NAME, - Key: key, - }), - ); - - log(`Deleted from S3: ${key}`); - } catch (err) { - log(`Failed to delete ${key}:`, err.message); - } - }); - } catch (error) { - log("Download error:", error.message); - - return res.status(404).json({ - success: false, - message: "Document not found", - }); - } -}; - -/** - * Cancel/delete a job - */ -exports.cancelJob = async (req, res) => { - try { - const { jobId } = req.params; - - if (!jobId) { - return res.status(400).json({ - success: false, - message: "jobId is required", - }); - } - - const job = await documentQueue.getJob(jobId); - - if (!job) { - return res.status(404).json({ - success: false, - message: "Job not found", - }); - } - - await job.remove(); - - return res.json({ - success: true, - message: "Job cancelled successfully", - jobId, - }); - } catch (error) { - log("Error cancelling job:", error.message); - return res.status(500).json({ - success: false, - message: error.message, - }); - } -}; +exports.getAvailableDocumentTypes = async (_req, res) => res.json({ success: true, data: Object.keys(registry) }); +exports.getSavedDocuments = async (req, res, next) => { try { const docs = await db.Document.findAll({ where: { owner_id: req.user.id }, attributes: { exclude: ["data"] }, order: [["createdAt", "DESC"]] }); return res.json({ success: true, data: docs.map(publicDoc) }); } catch (e) { return next(e); } }; +exports.getDocumentData = async (req, res, next) => { try { const doc = await db.Document.findByPk(req.params.docId); if (!doc) return res.status(404).json({ success: false, error: { code: "NOT_FOUND", message: "Document not found" } }); if (!allowed(req.user, doc)) return res.status(403).json({ success: false, error: { code: "FORBIDDEN", message: "Access denied" } }); return res.json({ success: true, data: { ...publicDoc(doc), documentData: doc.data } }); } catch (e) { return next(e); } }; +exports.generateReferenceNo = (_req, res) => res.status(410).json({ success: false, error: { code: "DEPRECATED", message: "Reference numbers are assigned during document creation" } }); +exports.generateDraftDocument = async (req, res, next) => { try { const parsed = requestSchema.safeParse({ ...req.body, documentType: req.body.documentType || "pdf" }); if (!parsed.success) return res.status(400).json({ success: false, error: { code: "VALIDATION_ERROR", message: "Invalid document request" } }); const key = normalize(parsed.data.document); if (!registry[key]) return res.status(400).json({ success: false, error: { code: "INVALID_DOCUMENT_TYPE", message: "Unsupported document type" } }); const doc = await db.Document.create({ doc_id: crypto.randomUUID(), reference_no: "N/A", doc_type: key.toUpperCase(), data: parsed.data.documentData || parsed.data.data || {}, status: "DRAFT", created_by: req.user.id, owner_type: "USER", owner_id: req.user.id }); return res.status(201).json({ success: true, data: publicDoc(doc) }); } catch (e) { return next(e); } }; +exports.generateDocument = async (req, res, next) => { try { const parsed = requestSchema.safeParse(req.body); if (!parsed.success) return res.status(400).json({ success: false, error: { code: "VALIDATION_ERROR", message: "Invalid document request", details: parsed.error.issues.map(i => ({ path: i.path.join("."), message: i.message })) } }); const key = normalize(parsed.data.document); const entry = registry[key]; if (!entry || (parsed.data.documentType === "excel" && !entry.excelBuilder)) return res.status(400).json({ success: false, error: { code: "INVALID_DOCUMENT_TYPE", message: "Document generator is unavailable" } }); const id = crypto.randomUUID(); const data = parsed.data.documentData || parsed.data.data || {}; const doc = await db.Document.create({ doc_id: id, reference_no: data.reference_no || "N/A", doc_type: key.toUpperCase(), data, status: "QUEUED", created_by: req.user.id, owner_type: "USER", owner_id: req.user.id, job_id: `document-${id}` }); try { await documentQueue.add("generate-document", { documentId: id, document: key, documentType: parsed.data.documentType, data }, { jobId: `document-${id}` }); } catch (e) { await doc.update({ status: "FAILED", failed_at: new Date(), failure_code: "QUEUE_FAILED" }); throw e; } return res.status(202).json({ success: true, data: publicDoc(doc) }); } catch (e) { return next(e); } }; +exports.getJobStatus = async (req, res, next) => { try { const doc = await db.Document.findOne({ where: { job_id: req.params.jobId } }); if (!doc) return res.status(404).json({ success: false, error: { code: "NOT_FOUND", message: "Job not found" } }); if (!allowed(req.user, doc)) return res.status(403).json({ success: false, error: { code: "FORBIDDEN", message: "Access denied" } }); return res.json({ success: true, data: publicDoc(doc) }); } catch (e) { return next(e); } }; +exports.downloadDocument = async (req, res, next) => { try { const doc = await db.Document.findByPk(req.params.docId, { include: [{ model: db.Upload, as: "storageUpload" }] }); if (!doc || doc.status !== "COMPLETED" || !doc.storageUpload) return res.status(404).json({ success: false, error: { code: "NOT_FOUND", message: "Document not available" } }); if (!allowed(req.user, doc)) return res.status(403).json({ success: false, error: { code: "FORBIDDEN", message: "Access denied" } }); const expiresIn = Number(process.env.S3_SIGNED_URL_TTL_SECONDS || 900); return res.json({ success: true, data: { url: await storage.createSignedDownloadUrl(doc.storageUpload.file_path, expiresIn), expiresIn } }); } catch (e) { return next(e); } }; +exports.cancelJob = async (req, res, next) => { try { const doc = await db.Document.findOne({ where: { job_id: req.params.jobId } }); if (!doc) return res.status(404).json({ success: false, error: { code: "NOT_FOUND", message: "Job not found" } }); if (!allowed(req.user, doc, "documents.delete")) return res.status(403).json({ success: false, error: { code: "FORBIDDEN", message: "Access denied" } }); const job = await documentQueue.getJob(doc.job_id); if (job) await job.remove(); await doc.update({ status: "FAILED", failed_at: new Date(), failure_code: "CANCELLED" }); return res.json({ success: true }); } catch (e) { return next(e); } }; diff --git a/app/controllers/notification.controller.js b/app/controllers/notification.controller.js index 61f0972..f3472d9 100644 --- a/app/controllers/notification.controller.js +++ b/app/controllers/notification.controller.js @@ -1,138 +1,20 @@ -/** - * 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/controllers/notification.controller.js - +const crypto = require("crypto"); const db = require("../models"); -const Notification = db.notification; -const UserNotification = db.userNotification; +const { Op } = require("sequelize"); -const log = require("../utils/consoleLog.utill").log; - -// Controller for managing notifications -exports.createNotification = async (req, res) => { +exports.createNotification = async (req, res, next) => { + const { notificationHeadline, notificationDescription, notificationType, userIds = [] } = req.body; + if (!notificationHeadline || !["USER", "ANNOUNCEMENT"].includes(notificationType)) return res.status(400).json({ success: false, error: { code: "VALIDATION_ERROR", message: "Valid headline and notificationType are required" } }); + if (notificationType === "USER" && (!Array.isArray(userIds) || !userIds.length)) return res.status(400).json({ success: false, error: { code: "VALIDATION_ERROR", message: "USER notifications require userIds" } }); + const transaction = await db.sequelize.transaction(); try { - const { notificationHeadline, notificationDescription, notificationType } = - req.body; - - const id = Date.now().toString(); // Generate a unique ID based on the current timestamp - const notification_id = `notif_${id}`; // Prefix the ID with "notif_" - - // Create a new notification - const newNotification = await Notification.create({ - notification_id, - notificationHeadline, - notificationDescription, - notificationType, - dateCreated: new Date(), - }); - - res.status(201).json({ - success: true, - message: "Notification created successfully", - notification: newNotification, - }); - } catch (error) { - log("Error creating notification:", error.message); - res.status(500).json({ - success: false, - message: "Internal server error", - error: error.message - }); - } + const notification = await db.notification.create({ notification_id: `notif_${crypto.randomUUID()}`, notificationHeadline, notificationDescription, notificationType, dateCreated: new Date() }, { transaction }); + const uniqueUsers = [...new Set(userIds)]; + if (uniqueUsers.length) await db.userNotification.bulkCreate(uniqueUsers.map((user_id) => ({ user_id, notification_id: notification.notification_id })), { transaction, ignoreDuplicates: true }); + await transaction.commit(); return res.status(201).json({ success: true, data: { notification, assignedUsers: uniqueUsers.length } }); + } catch (error) { await transaction.rollback(); return next(error); } }; - -// Mark a notification as read for a user -exports.markAsRead = async (req, res) => { - try { - const { userId, notificationId } = req.params; - - // Find the user notification entry - const userNotification = await UserNotification.findOne({ - where: { user_id: userId, notification_id: notificationId }, - }); - - if (!userNotification) { - return res.status(404).json({ success: false, error: "User notification not found" }); - } - - // Mark the notification as read - userNotification.isRead = true; - await userNotification.save(); - res.status(200).json({ - success: true, - message: "Notification marked as read", - }); - } - - catch (error) { - log("Error marking notification as read:", error.message); - res.status(500).json({ - success: false, - message: "Internal server error", - error: error.message - }); - } -}; - -// Get announcements for all users -exports.getAnnouncements = async (req, res) => { - try { - - // Fetch all announcements - const announcements = await Notification.findAll({ - where: { notificationType: "ANNOUNCEMENT", isActive: true }, - }); - - res.status(200).json({ - success: true, - announcements, - }); - - } catch (error) { - log("Error fetching announcements:", error.message); - res.status(500).json({ - success: false, - message: "Internal server error", - error: error.message - }); - } -} - -// Get all notifications for a user -exports.getUserNotifications = async (req, res) => { - try { - const { userId } = req.params; - - // Fetch notifications for the user - const notifications = await UserNotification.findAll({ - where: { user_id: userId, isRead: false }, - include: [ - { - model: Notification, - as: "notification", - }, - ], - }); - - res.status(200).json({ - success: true, - notifications, - }); - } - - catch (error) { - log("Error fetching user notifications:", error.message); - res.status(500).json({ - success: false, - message: "Internal server error", - error: error.message - }); - } -}; \ No newline at end of file +exports.getAnnouncements = async (_req, res, next) => { try { return res.json({ success: true, data: await db.notification.findAll({ where: { notificationType: "ANNOUNCEMENT", isActive: true }, order: [["createdAt", "DESC"]] }) }); } catch (e) { return next(e); } }; +exports.getMyNotifications = async (req, res, next) => { try { const rows = await db.userNotification.findAll({ where: { user_id: req.user.id }, include: [{ model: db.notification, as: "notification", where: { isActive: true } }], order: [["createdAt", "DESC"]] }); return res.json({ success: true, data: rows }); } catch (e) { return next(e); } }; +exports.markAsRead = async (req, res, next) => { try { const [count] = await db.userNotification.update({ isRead: true }, { where: { user_id: req.user.id, notification_id: req.params.notificationId } }); if (!count) return res.status(404).json({ success: false, error: { code: "NOT_FOUND", message: "Notification not found" } }); return res.json({ success: true }); } catch (e) { return next(e); } }; +exports.markAllAsRead = async (req, res, next) => { try { const [count] = await db.userNotification.update({ isRead: true }, { where: { user_id: req.user.id, isRead: false } }); return res.json({ success: true, data: { updated: count } }); } catch (e) { return next(e); } }; diff --git a/app/controllers/profile.controller.js b/app/controllers/profile.controller.js index b00457a..4960b9c 100644 --- a/app/controllers/profile.controller.js +++ b/app/controllers/profile.controller.js @@ -28,20 +28,21 @@ const { log } = require("../utils/consoleLog.utill"); // Get Profile avatar by user ID exports.getProfileAvatar = async (req, res) => { try { - const userId = req.params.userId; + const userId = req.params.userId || req.user.id; const profile = await Profile.findOne({ where: { user_id: userId } }); if (!profile) { return res.status(404).json({ error: "Profile not found" }); } const uploadRecord = await Upload.findByPk(profile.profilePicture_id); + if (!uploadRecord || uploadRecord.status !== "AVAILABLE" || (uploadRecord.visibility !== "PUBLIC" && uploadRecord.owner_id !== userId)) return res.status(404).json({ success: false, message: "Profile image not found" }); const fileUrl = await getSignedFileUrl(uploadRecord.file_path); res.json({ success: true, data: { - userId: profile.userId, + userId: profile.user_id, profilePictureUrl: fileUrl, }, }); @@ -58,7 +59,7 @@ exports.getProfileAvatar = async (req, res) => { // Get profile background image by user ID exports.getProfileBackgroundImage = async (req, res) => { try { - const userId = req.params.userId; + const userId = req.params.userId || req.user.id; const profile = await Profile.findOne({ where: { user_id: userId } }); if (!profile) { return res @@ -67,13 +68,14 @@ exports.getProfileBackgroundImage = async (req, res) => { } const uploadRecord = await Upload.findByPk(profile.backgroundImage_id); + if (!uploadRecord || uploadRecord.status !== "AVAILABLE" || (uploadRecord.visibility !== "PUBLIC" && uploadRecord.owner_id !== userId)) return res.status(404).json({ success: false, message: "Profile background not found" }); const fileUrl = await getSignedFileUrl(uploadRecord.file_path); res.json({ success: true, data: { - userId: profile.userId, + userId: profile.user_id, backgroundImageUrl: fileUrl, }, }); diff --git a/app/controllers/upload.controller.js b/app/controllers/upload.controller.js index ab93d12..eaebe9d 100644 --- a/app/controllers/upload.controller.js +++ b/app/controllers/upload.controller.js @@ -1,131 +1,44 @@ -/** - * 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/controllers/activity.controller.js - -const { uploadToS3, getSignedFileUrl } = require("../utils/s3Upload.utill"); -const { log } = require("../utils/consoleLog.utill"); - const db = require("../models"); -const Upload = db.Upload; -const Assets = db.Assets; +const storage = require("../services/storage/storage.service"); +const { validateUpload } = require("../services/storage/file-validation.service"); -// Constants -const MAX_IMAGE_SIZE = 3 * 1024 * 1024; // 3MB -const MAX_PDF_SIZE = 5 * 1024 * 1024; // 5MB +const canAccess = (user, upload) => upload.visibility === "PUBLIC" || upload.owner_id === user.id || user.accountType === "super_admin" || (user.accountType === "admin" && (user.permissions || []).includes("media.read")); -const isValidFileType = (mimetype) => { - return mimetype.startsWith("image/") || mimetype === "application/pdf"; -}; - -const isValidFileSize = (mimetype, size) => { - if (mimetype.startsWith("image/")) return size <= MAX_IMAGE_SIZE; - if (mimetype === "application/pdf") return size <= MAX_PDF_SIZE; - return false; -}; - -exports.uploadFile = async (req, res) => { +exports.uploadFile = async (req, res, next) => { + let objectKey; try { - // 1. Check file exists - if (!req.file) { - return res.status(400).json({ - success: false, - message: "No file uploaded", - }); - } - - const { mimetype, size, originalname } = req.file; - - // 2. Validate file type - if (!isValidFileType(mimetype)) { - return res.status(400).json({ - success: false, - message: "Only images and PDFs are allowed", - }); - } - - // 3. Validate file size - if (!isValidFileSize(mimetype, size)) { - return res.status(400).json({ - success: false, - message: mimetype.startsWith("image/") - ? "Image too large (max 1MB)" - : "PDF too large (max 5MB)", - }); - } - - log("File validation passed:", { - mimetype, - size, - originalname, - use_for: req.body.use_for, - }); - - // 4. Upload to S3 - const fileKey = await uploadToS3(req.file, req.body.use_for); - - const fileUrl = await getSignedFileUrl(fileKey); - - // 5. Save to DB - const uploadData = { - file_path: fileKey, - file_type: mimetype, - file_size: size, - original_name: originalname, - use_for: req.body.use_for || null, - uploaded_by: req.user?.id || null, - }; - - const newUpload = await Upload.create(uploadData); - - newUpload.dataValues.file_url = fileUrl; // Add URL to response - - // 6. Response - return res.status(201).json({ - success: true, - message: "File uploaded successfully", - data: newUpload, - }); - } catch (error) { - console.error("Upload Controller Error:", error); - - return res.status(500).json({ - success: false, - message: "Failed to upload file", - error: process.env.NODE_ENV === "development" ? error.message : undefined, - }); - } + const details = validateUpload(req.file); + const purpose = String(req.body.use_for || "generic").replace(/[^a-zA-Z0-9_-]/g, "_").slice(0, 60); + objectKey = storage.createObjectKey({ ownerId: req.user.id, mimeType: details.mimeType, purpose }); + await storage.uploadBuffer({ buffer: req.file.buffer, objectKey, mimeType: details.mimeType, checksum: details.checksum }); + let record; + try { + record = await db.Upload.create({ file_path: objectKey, file_type: details.mimeType, file_size: details.size, original_name: req.file.originalname, safe_name: details.safeName, checksum: details.checksum, use_for: purpose, uploaded_by: req.user.id, owner_type: "USER", owner_id: req.user.id, visibility: "PRIVATE", status: "AVAILABLE" }); + } catch (error) { await storage.deleteObject(objectKey).catch(() => undefined); throw error; } + const expiresIn = Number(process.env.S3_SIGNED_URL_TTL_SECONDS || 900); + return res.status(201).json({ success: true, data: { id: record.id, originalName: record.original_name, mimeType: record.file_type, size: record.file_size, purpose: record.use_for, status: record.status, url: await storage.createSignedDownloadUrl(objectKey, expiresIn), expiresIn } }); + } catch (error) { return next(error); } }; -// Get Uploaded File URL -exports.getFileUrl = async (req, res) => { +exports.getFileUrl = async (req, res, next) => { try { - const { id } = req.params; - - const uploadRecord = await Upload.findByPk(id); - - if (!uploadRecord) { - return res.status(404).json({ - success: false, - message: "File not found", - }); - } - - const fileUrl = await getSignedFileUrl(uploadRecord.file_path); - - return res.status(200).json({ - success: true, - message: "File URL retrieved successfully", - data: { - id: uploadRecord.id, - file_url: fileUrl, - }, - }); - } catch (error) {} + const upload = await db.Upload.findByPk(req.params.id); + if (!upload || upload.status === "DELETED") return res.status(404).json({ success: false, error: { code: "NOT_FOUND", message: "File not found" } }); + if (!canAccess(req.user, upload)) return res.status(403).json({ success: false, error: { code: "FORBIDDEN", message: "Access denied" } }); + if (upload.status !== "AVAILABLE") return res.status(409).json({ success: false, error: { code: "FILE_UNAVAILABLE", message: "File is not available" } }); + const expiresIn = Number(process.env.S3_SIGNED_URL_TTL_SECONDS || 900); + return res.json({ success: true, data: { id: upload.id, url: await storage.createSignedDownloadUrl(upload.file_path, expiresIn), expiresIn } }); + } catch (error) { return next(error); } +}; + +exports.deleteFile = async (req, res, next) => { + try { + const upload = await db.Upload.findByPk(req.params.id); + if (!upload || upload.status === "DELETED") return res.status(404).json({ success: false, error: { code: "NOT_FOUND", message: "File not found" } }); + const allowed = upload.owner_id === req.user.id || req.user.accountType === "super_admin" || (req.user.permissions || []).includes("media.delete"); + if (!allowed) return res.status(403).json({ success: false, error: { code: "FORBIDDEN", message: "Access denied" } }); + await upload.update({ status: "DELETED", deletedAt: new Date() }); + await storage.deleteObject(upload.file_path).catch(() => undefined); + return res.status(204).end(); + } catch (error) { return next(error); } }; diff --git a/app/logic/documents/registry.js b/app/logic/documents/registry.js index bd2be03..f76c65e 100644 --- a/app/logic/documents/registry.js +++ b/app/logic/documents/registry.js @@ -69,4 +69,8 @@ const registry = { console.log("📋 Registry initialized with keys:", Object.keys(registry)); -module.exports = registry; \ No newline at end of file +for (const [key, entry] of Object.entries(registry)) { + if (typeof entry.pdfTemplate !== "function") delete registry[key]; +} + +module.exports = registry; diff --git a/app/middleware/docsSession.middleware.js b/app/middleware/docsSession.middleware.js index b7d23e8..5e3a2a7 100644 --- a/app/middleware/docsSession.middleware.js +++ b/app/middleware/docsSession.middleware.js @@ -1,21 +1,6 @@ -/** - * 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/middleware/docsSession.middleware.js - -module.exports = (req, res, next) => { - if (req.cookies && req.cookies.docsAuth === "true") { - return next(); - } - - return res.status(401).json({ - success: false, - message: "Unauthorized: Please login to access docs", - }); -}; \ No newline at end of file +const { authenticate } = require("./auth.middleware"); +const { ACCOUNT_TYPES } = require("../constants/accountTypes"); +module.exports = (req, res, next) => authenticate(req, res, () => { + if ([ACCOUNT_TYPES.ADMIN, ACCOUNT_TYPES.SUPER_ADMIN].includes(req.user.accountType)) return next(); + return res.status(403).json({ success: false, error: { code: "FORBIDDEN", message: "Administrative access required" } }); +}); diff --git a/app/middleware/upload.middleware.js b/app/middleware/upload.middleware.js index b4dd001..00b22ac 100644 --- a/app/middleware/upload.middleware.js +++ b/app/middleware/upload.middleware.js @@ -16,23 +16,12 @@ const storage = multer.memoryStorage(); // File filter (only images + PDFs) const fileFilter = (req, file, cb) => { - const allowedTypes = [ - "image/", - "application/pdf", - "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet", // .xlsx - "application/vnd.ms-excel" // .xls - ]; - - const isValid = - file.mimetype.startsWith("image/") || - file.mimetype === "application/pdf" || - file.mimetype === "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet" || - file.mimetype === "application/vnd.ms-excel"; + const isValid = ["image/jpeg", "image/png", "image/webp", "application/pdf", "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet"].includes(file.mimetype); if (isValid) { cb(null, true); } else { - cb(new Error("Only images, PDFs, and Excel files are allowed"), false); + cb(new Error("Only JPEG, PNG, WebP, PDF, and XLSX files are allowed"), false); } }; @@ -41,7 +30,7 @@ const upload = multer({ storage, fileFilter, limits: { - fileSize: 5 * 1024 * 1024, // max 5MB (global limit) + fileSize: Number(process.env.S3_MAX_UPLOAD_BYTES || 5 * 1024 * 1024), }, }); @@ -61,4 +50,4 @@ module.exports = { uploadSingle, uploadMultiple, uploadFields, -}; \ No newline at end of file +}; diff --git a/app/models/activities/userActivities.model.js b/app/models/activities/userActivities.model.js index 8560b29..cfc0fc5 100644 --- a/app/models/activities/userActivities.model.js +++ b/app/models/activities/userActivities.model.js @@ -18,9 +18,10 @@ module.exports = (sequelize, DataTypes) => { autoIncrement: true, primaryKey: true, }, + event_id: { type: DataTypes.STRING, allowNull: true, unique: true }, user_id: { type: DataTypes.STRING, - allowNull: false, + allowNull: true, }, username: { type: DataTypes.STRING, @@ -45,13 +46,17 @@ module.exports = (sequelize, DataTypes) => { activity_date:{ type: DataTypes.DATEONLY, allowNull: false, - } + }, + target_type: DataTypes.STRING, target_id: DataTypes.STRING, + ip_address: DataTypes.STRING, user_agent: DataTypes.STRING, request_id: DataTypes.STRING, + metadata: DataTypes.JSON, occurred_at: { type: DataTypes.DATE, allowNull: false, defaultValue: DataTypes.NOW }, }, { tableName: "UserActivity", timestamps: true, }, ); + UserActivity.associate = (models) => UserActivity.belongsTo(models.User, { foreignKey: "user_id", as: "actor", constraints: false }); return UserActivity; }; diff --git a/app/models/document/document.model.js b/app/models/document/document.model.js index 0bd05ca..2854b43 100644 --- a/app/models/document/document.model.js +++ b/app/models/document/document.model.js @@ -14,7 +14,7 @@ module.exports = (sequelize, DataTypes) => { "Document", { doc_id:{ - type: DataTypes.STRING, + type: DataTypes.ENUM("DRAFT", "QUEUED", "PROCESSING", "COMPLETED", "FAILED"), primaryKey: true, }, reference_no: { @@ -34,13 +34,21 @@ module.exports = (sequelize, DataTypes) => { type: DataTypes.STRING, allowNull: false, defaultValue: "DRAFT", - } + }, + created_by: { type: DataTypes.STRING, allowNull: false }, + owner_type: { type: DataTypes.STRING, allowNull: false, defaultValue: "USER" }, + owner_id: { type: DataTypes.STRING, allowNull: false }, + storage_upload_id: { type: DataTypes.INTEGER, allowNull: true }, + job_id: { type: DataTypes.STRING, allowNull: true, unique: true }, + generated_at: DataTypes.DATE, failed_at: DataTypes.DATE, failure_code: DataTypes.STRING, + version: { type: DataTypes.INTEGER, allowNull: false, defaultValue: 1 }, }, { tableName: "Document", timestamps: true, }, ); + document.associate = (models) => { document.belongsTo(models.User, { foreignKey: "created_by", as: "creator", constraints: false }); document.belongsTo(models.Upload, { foreignKey: "storage_upload_id", as: "storageUpload", constraints: false }); }; return document; }; diff --git a/app/models/document/documentType.model.js b/app/models/document/documentType.model.js index 504eec2..ef76a3a 100644 --- a/app/models/document/documentType.model.js +++ b/app/models/document/documentType.model.js @@ -20,8 +20,8 @@ module.exports = (sequelize, DataTypes) => { }, doc_type_name: { type: DataTypes.STRING, - allowNull: true, - defaultValue: "N/A" + allowNull: false, + unique: true, }, description: { type: DataTypes.JSON, diff --git a/app/models/index.js b/app/models/index.js index 2b6b999..9f70090 100644 --- a/app/models/index.js +++ b/app/models/index.js @@ -64,6 +64,7 @@ db.referenceNumber = require("./referenceNumbers/referenceNumber.model")(sequeli // Notifications db.notification = require("./notification/notification.model")(sequelize, DataTypes); db.userNotification = require("./notification/userNotification.model")(sequelize, DataTypes); +db.NotificationDelivery = require("./notification/notificationDelivery.model")(sequelize, DataTypes); /* Associations */ Object.keys(db).forEach(model => { diff --git a/app/models/notification/notificationDelivery.model.js b/app/models/notification/notificationDelivery.model.js new file mode 100644 index 0000000..4b07e64 --- /dev/null +++ b/app/models/notification/notificationDelivery.model.js @@ -0,0 +1,11 @@ +module.exports = (sequelize, DataTypes) => sequelize.define("NotificationDelivery", { + id: { type: DataTypes.STRING, primaryKey: true }, + notificationId: { type: DataTypes.STRING, allowNull: true }, + userId: { type: DataTypes.STRING, allowNull: true }, + channel: { type: DataTypes.ENUM("EMAIL", "PUSH"), allowNull: false }, + recipient: { type: DataTypes.STRING, allowNull: false }, + templateKey: { type: DataTypes.STRING, allowNull: false }, + status: { type: DataTypes.ENUM("QUEUED", "PROCESSING", "SENT", "FAILED"), allowNull: false, defaultValue: "QUEUED" }, + attemptCount: { type: DataTypes.INTEGER, allowNull: false, defaultValue: 0 }, + lastAttemptAt: DataTypes.DATE, sentAt: DataTypes.DATE, failedAt: DataTypes.DATE, failureCode: DataTypes.STRING, +}, { tableName: "notification_deliveries", timestamps: true }); diff --git a/app/models/notification/userNotification.model.js b/app/models/notification/userNotification.model.js index ce3d154..b50ec64 100644 --- a/app/models/notification/userNotification.model.js +++ b/app/models/notification/userNotification.model.js @@ -41,11 +41,13 @@ module.exports = (sequelize, DataTypes) => { userNotification.associate = (models) => { userNotification.belongsTo(models.User, { foreignKey: "user_id", + as: "user", constraints: false, }); userNotification.belongsTo(models.notification, { foreignKey: "notification_id", + as: "notification", constraints: false, }); }; diff --git a/app/models/upload/upload.model.js b/app/models/upload/upload.model.js index a345a29..23b904c 100644 --- a/app/models/upload/upload.model.js +++ b/app/models/upload/upload.model.js @@ -43,12 +43,21 @@ module.exports = (sequelize, DataTypes) => { type: DataTypes.STRING, allowNull: false, }, + storage_provider: { type: DataTypes.STRING, allowNull: false, defaultValue: "S3" }, + safe_name: { type: DataTypes.STRING, allowNull: true }, + checksum: { type: DataTypes.STRING(64), allowNull: true }, + owner_type: { type: DataTypes.STRING, allowNull: false, defaultValue: "USER" }, + owner_id: { type: DataTypes.STRING, allowNull: false }, + visibility: { type: DataTypes.ENUM("PRIVATE", "PUBLIC"), allowNull: false, defaultValue: "PRIVATE" }, + status: { type: DataTypes.ENUM("PENDING", "AVAILABLE", "FAILED", "DELETED"), allowNull: false, defaultValue: "PENDING" }, + deletedAt: { type: DataTypes.DATE, allowNull: true }, }, { tableName: "uploads", timestamps: true, }, ); + upload.associate = (models) => upload.belongsTo(models.User, { foreignKey: "uploaded_by", as: "uploader", constraints: false }); return upload; }; diff --git a/app/queues/activity.queue.js b/app/queues/activity.queue.js index eea6434..f8828c5 100644 --- a/app/queues/activity.queue.js +++ b/app/queues/activity.queue.js @@ -15,6 +15,7 @@ const connection = require("../config/redisClient"); const activityQueue = new Queue("activity-queue", { connection, + defaultJobOptions: { attempts: 3, backoff: { type: "exponential", delay: 1000 }, removeOnComplete: { age: 86400, count: 5000 }, removeOnFail: { age: 604800, count: 5000 } }, }); -module.exports = activityQueue; \ No newline at end of file +module.exports = activityQueue; diff --git a/app/queues/document.queue.js b/app/queues/document.queue.js index 3195261..e39293d 100644 --- a/app/queues/document.queue.js +++ b/app/queues/document.queue.js @@ -27,7 +27,7 @@ const documentQueue = new Queue("document-generation", { count: 1000, }, - removeOnFail: false, + removeOnFail: { age: 604800, count: 5000 }, }, }); diff --git a/app/queues/email.queue.js b/app/queues/email.queue.js new file mode 100644 index 0000000..d0b27d4 --- /dev/null +++ b/app/queues/email.queue.js @@ -0,0 +1,3 @@ +const { Queue } = require("bullmq"); +const connection = require("../config/redisClient"); +module.exports = new Queue("email-delivery", { connection, defaultJobOptions: { attempts: 5, backoff: { type: "exponential", delay: 2000 }, removeOnComplete: { age: 3600, count: 1000 }, removeOnFail: { age: 604800, count: 5000 } } }); diff --git a/app/queues/log.queue.js b/app/queues/log.queue.js index fe89ffe..d5b3b58 100644 --- a/app/queues/log.queue.js +++ b/app/queues/log.queue.js @@ -12,6 +12,6 @@ const { Queue } = require("bullmq"); const connection = require("../config/redisClient"); -const logQueue = new Queue("logQueue", { connection }); +const logQueue = new Queue("logQueue", { connection, defaultJobOptions: { attempts: 2, backoff: { type: "fixed", delay: 1000 }, removeOnComplete: { age: 3600, count: 5000 }, removeOnFail: { age: 86400, count: 5000 } } }); -module.exports = logQueue; \ No newline at end of file +module.exports = logQueue; diff --git a/app/routes/activity.routes.js b/app/routes/activity.routes.js index 2855cda..b962ac0 100644 --- a/app/routes/activity.routes.js +++ b/app/routes/activity.routes.js @@ -24,7 +24,7 @@ const PERMISSIONS = require("../constants/permissions"); router.get( "/", authenticate, - authorizedAccountType(["admin"]), + checkPermission("audit.read", { custom: true }), activityController.getUserActivities, ); @@ -32,7 +32,7 @@ router.get( router.get( "/user/:userId", authenticate, - authorizedAccountType(["admin"]), + checkPermission("audit.read", { custom: true }), activityController.getActivitiesByUserId, ); diff --git a/app/routes/docs.routes.js b/app/routes/docs.routes.js index 261bed2..b00234c 100644 --- a/app/routes/docs.routes.js +++ b/app/routes/docs.routes.js @@ -48,31 +48,8 @@ router.get("/view/:file", docsSession, (req, res) => { }); // Docs login (simple) -router.post("/login", sensitiveLimiter, (req, res) => { - const { username, password } = req.body; +router.post("/login", sensitiveLimiter, (_req, res) => res.status(410).json({ success: false, error: { code: "DEPRECATED", message: "Use an authenticated ADMIN or SUPER_ADMIN session" } })); - if ( - username === process.env.DOCS_USER && - password === process.env.DOCS_PASS - ) { - res.cookie("docsAuth", "true", { - httpOnly: true, - sameSite: "lax", - secure: process.env.NODE_ENV === "production", - }); - - return res.json({ success: true }); - } - - return res.status(401).json({ - success: false, - message: "Invalid credentials", - }); -}); - -router.post("/logout", (req, res) => { - res.clearCookie("docsAuth"); - res.json({ success: true }); -}); +router.post("/logout", (_req, res) => res.status(410).json({ success: false, error: { code: "DEPRECATED", message: "Use the authentication logout endpoint" } })); module.exports = router; diff --git a/app/routes/document.routes.js b/app/routes/document.routes.js index ce1cb1f..8279666 100644 --- a/app/routes/document.routes.js +++ b/app/routes/document.routes.js @@ -1,94 +1,12 @@ -/** - * 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/routes/document.routes.js - -const express = require("express"); -const router = express.Router(); -const generateDocumentController = require("../controllers/generateDocument.controller"); -const { - authenticate, -} = require("../middleware/auth.middleware"); - -const { authorizedAccountType, checkPermission } = require("../middleware/permission.middleware"); -const PERMISSIONS = require("../constants/permissions"); - - -// Get available document types -router.get( - "/types", - authenticate, - authorizedAccountType(["admin", "management", "team_head", "user"]), - generateDocumentController.getAvailableDocumentTypes -); - -// Get saved documents -router.post( - "/saved", - authenticate, - authorizedAccountType(["admin", "management", "team_head", "user"]), - generateDocumentController.getSavedDocuments -); - -// Generate Document (async - returns jobId immediately) -router.post( - "/generate", - authenticate, - authorizedAccountType(["admin", "management", "team_head", "user"]), - generateDocumentController.generateDocument -); - -router.post( - "/draft", - authenticate, - authorizedAccountType(["admin", "management", "team_head", "user"]), - generateDocumentController.generateDraftDocument -); - -// Generate Reference Number -router.get( - "/reference-number/:documentType", - authenticate, - authorizedAccountType(["admin", "management", "team_head", "user"]), - generateDocumentController.generateReferenceNo -); - -// Get Job Status -router.get( - "/job/:jobId/status", - authenticate, - authorizedAccountType(["admin", "management", "team_head", "user"]), - generateDocumentController.getJobStatus -); - -// Download Document -router.get( - "/download/:uuid", - authenticate, - authorizedAccountType(["admin", "management", "team_head", "user"]), - generateDocumentController.downloadDocument -); - -// Cancel Job -router.delete( - "/job/:jobId", - authenticate, - authorizedAccountType(["admin", "management", "team_head", "user"]), - generateDocumentController.cancelJob -); - -// Get document details -router.get( - "/:docId", - authenticate, - authorizedAccountType(["admin", "management", "team_head", "user"]), - generateDocumentController.getDocumentData -); - -module.exports = router; \ No newline at end of file +const router = require("express").Router(); +const c = require("../controllers/generateDocument.controller"); +const { authenticate } = require("../middleware/auth.middleware"); +router.use(authenticate); +router.get("/types", c.getAvailableDocumentTypes); +router.get("/saved", c.getSavedDocuments); router.post("/saved", c.getSavedDocuments); +router.post("/generate", c.generateDocument); router.post("/draft", c.generateDraftDocument); +router.get("/reference-number/:documentType", c.generateReferenceNo); +router.get("/jobs/:jobId", c.getJobStatus); router.get("/job/:jobId/status", c.getJobStatus); +router.get("/:docId/download", c.downloadDocument); router.get("/download/:docId", c.downloadDocument); +router.delete("/job/:jobId", c.cancelJob); router.get("/:docId", c.getDocumentData); +module.exports = router; diff --git a/app/routes/notification.routes.js b/app/routes/notification.routes.js index 5b328ca..5840151 100644 --- a/app/routes/notification.routes.js +++ b/app/routes/notification.routes.js @@ -1,69 +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. - */ - -// app/routes/notification.routes.js - -const express = require("express"); -const router = express.Router(); - -const notificationController = require("../controllers/notification.controller"); - -const { - authenticate, -} = require("../middleware/auth.middleware"); - -const { - authorizedAccountType, - checkPermission, -} = require("../middleware/permission.middleware"); - -const PERMISSIONS = require("../constants/permissions"); - - -// Create notification -router.post( - "/", - authenticate, - authorizedAccountType(["admin", "management"]), - // checkPermission(PERMISSIONS.NOTIFICATION_CREATE), - notificationController.createNotification, -); - - -// Get all announcements -router.get( - "/announcements", - authenticate, - authorizedAccountType(["admin", "management", "team_head", "user"]), - // checkPermission(PERMISSIONS.NOTIFICATION_VIEW), - notificationController.getAnnouncements, -); - - -// Get notifications for a user -router.get( - "/user/:userId", - authenticate, - authorizedAccountType(["admin", "management", "team_head", "user"]), - // checkPermission(PERMISSIONS.NOTIFICATION_VIEW), - notificationController.getUserNotifications, -); - - -// Mark notification as read -router.patch( - "/user/:userId/notification/:notificationId/read", - authenticate, - authorizedAccountType(["admin", "management", "team_head", "user"]), - // checkPermission(PERMISSIONS.NOTIFICATION_READ), - notificationController.markAsRead, -); - - -module.exports = router; \ No newline at end of file +const router = require("express").Router(); +const controller = require("../controllers/notification.controller"); +const { authenticate } = require("../middleware/auth.middleware"); +const { checkPermission } = require("../middleware/permission.middleware"); +router.use(authenticate); +router.post("/", checkPermission("notifications.manage", { custom: true }), controller.createNotification); +router.get("/announcements", controller.getAnnouncements); +router.get("/me", controller.getMyNotifications); +router.patch("/read-all", controller.markAllAsRead); +router.patch("/:notificationId/read", controller.markAsRead); +module.exports = router; diff --git a/app/routes/profile.routes.js b/app/routes/profile.routes.js index 8e92643..ad76262 100644 --- a/app/routes/profile.routes.js +++ b/app/routes/profile.routes.js @@ -36,6 +36,16 @@ router.post( authController.changePassword, ); +router.get( + "/me/avatar", + authenticate, + profileController.getProfileAvatar, +); +router.get( + "/me/background", + authenticate, + profileController.getProfileBackgroundImage, +); router.get( "/avatar/:userId", authenticate, diff --git a/app/routes/upload.routes.js b/app/routes/upload.routes.js index 00c0c7f..7e1bb0e 100644 --- a/app/routes/upload.routes.js +++ b/app/routes/upload.routes.js @@ -24,7 +24,6 @@ const PERMISSIONS = require("../constants/permissions"); router.post( "/", authenticate, - authorizedAccountType(["admin", "staff"]), uploadSingle("file"), uploadController.uploadFile, ); @@ -33,8 +32,8 @@ router.post( router.get( "/signed-url/:id", authenticate, - authorizedAccountType(["admin", "staff"]), uploadController.getFileUrl, ); +router.delete("/:id", authenticate, uploadController.deleteFile); -module.exports = router; \ No newline at end of file +module.exports = router; diff --git a/app/services/activity.service.js b/app/services/activity.service.js index 110d468..0e03da3 100644 --- a/app/services/activity.service.js +++ b/app/services/activity.service.js @@ -15,21 +15,33 @@ const logActivity = async ({ user, description, type = "ACTION", - module + module, + eventId, + targetType, + targetId, + requestId, + ipAddress, + userAgent, + metadata, }) => { try { + const safeMetadata = metadata && JSON.parse(JSON.stringify(metadata, (key, value) => /password|otp|token|authorization|cookie/i.test(key) ? "[REDACTED]" : value)); + const id = eventId || require("crypto").randomUUID(); await activityQueue.add("log-activity", { - user_id: user.id || user.user_id, - username: user.firstName, + event_id: id, + user_id: user?.id || user?.user_id || null, + username: user?.firstName || "System", activity_description: description, module: module, activity_type: type, activity_time: new Date(), - activity_date: new Date().toISOString().split("T")[0] - }); + activity_date: new Date().toISOString().split("T")[0], + target_type: targetType, target_id: targetId, request_id: requestId, + ip_address: ipAddress, user_agent: userAgent, metadata: safeMetadata, occurred_at: new Date(), + }, { jobId: `activity-${id}` }); } catch (err) { console.error("Queue push error:", err.message); } }; -module.exports = { logActivity }; \ No newline at end of file +module.exports = { logActivity }; diff --git a/app/services/auth/email.service.js b/app/services/auth/email.service.js index a77a866..5fd4865 100644 --- a/app/services/auth/email.service.js +++ b/app/services/auth/email.service.js @@ -1,6 +1,6 @@ -const { sendMail } = require("../../utils/mail.util"); +const emailService = require("../email/email.service"); -const sendLoginOtp = (user, otp) => sendMail({ to: user.email, subject: "Your ZUMRI login code", templateName: "otp", templateVars: { firstName: user.firstName, otp }, text: `Your ZUMRI login code is ${otp}.` }); -const sendPasswordChanged = (user) => sendMail({ to: user.email, subject: "Your ZUMRI password was changed", templateName: "passwordChanged", templateVars: { customer_name: user.firstName, changed_at: new Date().toLocaleString() }, text: "Your ZUMRI password was changed. Contact support if this was not you." }); +const sendLoginOtp = (user, otp) => emailService.send({ templateKey: "otp", recipient: user.email, userId: user.id, variables: { firstName: user.firstName, otp } }); +const sendPasswordChanged = (user) => emailService.send({ templateKey: "passwordChanged", recipient: user.email, userId: user.id, variables: { customer_name: user.firstName, changed_at: new Date().toLocaleString() } }); module.exports = { sendLoginOtp, sendPasswordChanged }; diff --git a/app/services/email/email.service.js b/app/services/email/email.service.js new file mode 100644 index 0000000..c68d78e --- /dev/null +++ b/app/services/email/email.service.js @@ -0,0 +1,15 @@ +const crypto = require("crypto"); +const emailQueue = require("../../queues/email.queue"); +const db = require("../../models"); +const templates = { + otp: { subject: "Your ZUMRI login code" }, passwordChanged: { subject: "Your ZUMRI password was changed" }, + passwordReset: { subject: "Reset your ZUMRI password" }, emailVerification: { subject: "Verify your ZUMRI email" }, welcome: { subject: "Welcome to ZUMRI" }, +}; +const send = async ({ templateKey, recipient, locale = "en", variables = {}, correlationId, notificationId, userId, eventId }) => { + if (!templates[templateKey]) throw new Error("Unknown email template"); + const id = crypto.randomUUID(); + await db.NotificationDelivery.create({ id, notificationId, userId, channel: "EMAIL", recipient, templateKey, status: "QUEUED" }); + await emailQueue.add("send-email", { deliveryId: id, templateKey, recipient, locale, variables, correlationId }, { jobId: `email-${eventId || id}` }); + return { deliveryId: id }; +}; +module.exports = { send, templates }; diff --git a/app/services/notification/notification-policy.service.js b/app/services/notification/notification-policy.service.js new file mode 100644 index 0000000..8ca3d2d --- /dev/null +++ b/app/services/notification/notification-policy.service.js @@ -0,0 +1,7 @@ +const MANDATORY_SECURITY_EVENTS = new Set(["LOGIN_OTP", "PASSWORD_RESET", "PASSWORD_CHANGED", "EMAIL_VERIFICATION"]); +const channelAllowed = ({ eventType, channel, profile }) => { + if (MANDATORY_SECURITY_EVENTS.has(eventType)) return true; + if (channel === "IN_APP") return true; + return profile?.notificationsEnabled !== false; +}; +module.exports = { MANDATORY_SECURITY_EVENTS, channelAllowed }; diff --git a/app/services/storage/file-validation.service.js b/app/services/storage/file-validation.service.js new file mode 100644 index 0000000..f0d04d6 --- /dev/null +++ b/app/services/storage/file-validation.service.js @@ -0,0 +1,19 @@ +const crypto = require("crypto"); +const { extensionFor, sanitizeFilename } = require("./storage.service"); +const detectMimeType = (buffer) => { + if (!Buffer.isBuffer(buffer) || !buffer.length) return null; + if (buffer.subarray(0, 3).equals(Buffer.from([0xff, 0xd8, 0xff]))) return "image/jpeg"; + if (buffer.subarray(0, 8).equals(Buffer.from([0x89,0x50,0x4e,0x47,0x0d,0x0a,0x1a,0x0a]))) return "image/png"; + if (buffer.length >= 12 && buffer.toString("ascii", 0, 4) === "RIFF" && buffer.toString("ascii", 8, 12) === "WEBP") return "image/webp"; + if (buffer.subarray(0, 5).toString("ascii") === "%PDF-") return "application/pdf"; + if (buffer.subarray(0, 4).equals(Buffer.from([0x50,0x4b,0x03,0x04]))) return "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet"; + return null; +}; +const validateUpload = (file) => { + if (!file?.buffer?.length) throw Object.assign(new Error("Empty files are not allowed"), { status: 400 }); + if (file.buffer.length > Number(process.env.S3_MAX_UPLOAD_BYTES || 5242880)) throw Object.assign(new Error("File exceeds the upload size limit"), { status: 413 }); + const mimeType = detectMimeType(file.buffer); + if (!mimeType || mimeType !== file.mimetype || !extensionFor(mimeType)) throw Object.assign(new Error("File content does not match an allowed type"), { status: 400 }); + return { mimeType, size: file.buffer.length, checksum: crypto.createHash("sha256").update(file.buffer).digest("hex"), safeName: `${sanitizeFilename(file.originalname).replace(/\.[^.]+$/, "")}${extensionFor(mimeType)}` }; +}; +module.exports = { detectMimeType, validateUpload }; diff --git a/app/services/storage/storage.service.js b/app/services/storage/storage.service.js new file mode 100644 index 0000000..a325ed5 --- /dev/null +++ b/app/services/storage/storage.service.js @@ -0,0 +1,19 @@ +const crypto = require("crypto"); +const path = require("path"); +const { PutObjectCommand, DeleteObjectCommand, HeadObjectCommand, GetObjectCommand } = require("@aws-sdk/client-s3"); +const { getSignedUrl } = require("@aws-sdk/s3-request-presigner"); +const s3 = require("../../config/s3.config"); +const bucket = () => process.env.AWS_S3_BUCKET_NAME; +const extensionFor = (mime) => ({ "image/jpeg": ".jpg", "image/png": ".png", "image/webp": ".webp", "application/pdf": ".pdf", "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet": ".xlsx" }[mime]); +const sanitizeFilename = (name = "file") => path.basename(name).replace(/[^a-zA-Z0-9._-]/g, "_").slice(0, 180); +const createObjectKey = ({ ownerId = "system", mimeType, purpose = "upload", now = new Date(), id = crypto.randomUUID() }) => { + const ext = extensionFor(mimeType); if (!ext) throw new Error("Unsupported file type"); + const root = purpose === "document" ? "documents" : "uploads"; + const owner = String(ownerId).replace(/[^a-zA-Z0-9_-]/g, "_").slice(0, 80); + return `${root}/${owner}/${now.getUTCFullYear()}/${String(now.getUTCMonth() + 1).padStart(2, "0")}/${id}${ext}`; +}; +const uploadBuffer = async ({ buffer, objectKey, mimeType, checksum }) => { await s3.send(new PutObjectCommand({ Bucket: bucket(), Key: objectKey, Body: buffer, ContentType: mimeType, ChecksumSHA256: checksum ? Buffer.from(checksum, "hex").toString("base64") : undefined })); return objectKey; }; +const deleteObject = (objectKey) => s3.send(new DeleteObjectCommand({ Bucket: bucket(), Key: objectKey })); +const objectExists = async (objectKey) => { try { await s3.send(new HeadObjectCommand({ Bucket: bucket(), Key: objectKey })); return true; } catch (error) { if (error.name === "NotFound" || error.$metadata?.httpStatusCode === 404) return false; throw error; } }; +const createSignedDownloadUrl = (objectKey, expiresIn = Number(process.env.S3_SIGNED_URL_TTL_SECONDS || 900)) => getSignedUrl(s3, new GetObjectCommand({ Bucket: bucket(), Key: objectKey }), { expiresIn }); +module.exports = { uploadBuffer, deleteObject, objectExists, createSignedDownloadUrl, createObjectKey, extensionFor, sanitizeFilename }; diff --git a/app/utils/consoleLog.utill.js b/app/utils/consoleLog.utill.js index 33193cf..adccd2c 100644 --- a/app/utils/consoleLog.utill.js +++ b/app/utils/consoleLog.utill.js @@ -23,11 +23,11 @@ const log = async (...args) => { const message = args .map((arg) => { if (arg instanceof Error) { - return `${arg.message}\n${arg.stack}`; + return arg.message; } if (typeof arg === "object") { - return JSON.stringify(arg); + return JSON.stringify(arg, (key, value) => /password|otp|token|authorization|cookie|secret/i.test(key) ? "[REDACTED]" : value); } return String(arg); @@ -35,7 +35,7 @@ const log = async (...args) => { .join(" "); await logQueue.add("log", { - message, + level: "info", service: "zumri-api", message, timestamp: new Date().toISOString(), }); } catch (err) { diff --git a/app/utils/documentJob.util.js b/app/utils/documentJob.util.js index 1806c4f..67b06d9 100644 --- a/app/utils/documentJob.util.js +++ b/app/utils/documentJob.util.js @@ -9,7 +9,7 @@ // app/utils/documentJob.util.js -const documentQueue = require("../queues/pdf.queue"); +const documentQueue = require("../queues/document.queue"); /** * Queue a document generation job diff --git a/app/utils/emailVerification.util.js b/app/utils/emailVerification.util.js index c5a6a7f..af3d695 100644 --- a/app/utils/emailVerification.util.js +++ b/app/utils/emailVerification.util.js @@ -1,7 +1,7 @@ // app/utils/emailVerification.util.js const crypto = require("crypto"); -const { sendMail } = require("../utils/mail.util"); +const emailService = require("../services/email/email.service"); const redis = require("../config/redisClient"); const EMAIL_VERIFICATION_TTL = @@ -67,19 +67,12 @@ const sendVerificationEmail = async (email, firstName, verificationToken) => { const confirmationLink = `${process.env.FRONTEND_URL}/verify-email?token=${verificationToken}`; // Send email - await sendMail({ - to: email, - - subject: "Confirm Your ZUMRI Account", - - templateName: "emailVerification", - - templateVars: { + await emailService.send({ + recipient: email, templateKey: "emailVerification", + variables: { customer_name: firstName, confirmation_link: confirmationLink, - }, - - text: `Hello ${firstName}, please verify your ZUMRI account using this link: ${confirmationLink}`, + }, eventId: `verify-${hashEmailVerificationToken(verificationToken)}`, }); }; diff --git a/app/utils/mail.util.js b/app/utils/mail.util.js index aa94ee7..ba64def 100644 --- a/app/utils/mail.util.js +++ b/app/utils/mail.util.js @@ -32,12 +32,14 @@ if (process.env.NODE_ENV !== "test" && process.env.ENABLE_MAIL === "true") { // Utility to load template and replace placeholders const loadTemplate = (templateName, variables = {}) => { + if (!/^[a-zA-Z0-9_-]+$/.test(templateName)) throw new Error("Invalid email template"); const templatePath = path.join(__dirname, "../templates/emails", `${templateName}.html`); let template = fs.readFileSync(templatePath, "utf-8"); Object.keys(variables).forEach((key) => { const regex = new RegExp(`{{${key}}}`, "g"); - template = template.replace(regex, variables[key]); + const escaped = String(variables[key] ?? "").replace(/[&<>"']/g, (char) => ({ "&": "&", "<": "<", ">": ">", '"': """, "'": "'" })[char]); + template = template.replace(regex, escaped); }); return template; @@ -90,5 +92,5 @@ const sendMail = async (options) => { }; -module.exports = { sendMail }; +module.exports = { sendMail, loadTemplate }; diff --git a/app/utils/passwordReset.utill.js b/app/utils/passwordReset.utill.js index ecc7864..4ac211b 100644 --- a/app/utils/passwordReset.utill.js +++ b/app/utils/passwordReset.utill.js @@ -11,7 +11,7 @@ const crypto = require("crypto"); const redis = require("../config/redisClient"); -const { sendMail } = require("./mail.util"); +const emailService = require("../services/email/email.service"); const PASSWORD_RESET_TTL = Number(process.env.PASSWORD_RESET_TTL_SECONDS) || 900; @@ -81,23 +81,13 @@ const sendPasswordResetEmail = `${process.env.FRONTEND_URL}/reset-password?token=${resetToken}`; - await sendMail({ - to: email, - - subject: - "Reset your ZUMRI password", - - templateName: - "passwordReset", - - templateVars: { + await emailService.send({ + recipient: email, templateKey: "passwordReset", + variables: { firstName: firstName, resetLink: resetLink, expiryTime: "15 minutes", - }, - - text: - `Hi ${firstName}, use this link within 15 minutes to reset your password: ${resetLink}`, + }, eventId: `reset-${hashPasswordResetToken(resetToken)}`, }); }; @@ -109,25 +99,15 @@ const sendPasswordResetEmail = const changedAt = new Date().toLocaleString(); - await sendMail({ - to: email, - - subject: - "Your ZUMRI password was changed", - - templateName: - "passwordChanged", - - templateVars: { + await emailService.send({ + recipient: email, templateKey: "passwordChanged", + variables: { customer_name: firstName, changed_at: changedAt, }, - - text: - `Hi ${firstName}, your password was changed at ${changedAt}. If this was not you, contact support immediately.`, }); }; diff --git a/app/utils/s3Upload.utill.js b/app/utils/s3Upload.utill.js index c2eacae..574499c 100644 --- a/app/utils/s3Upload.utill.js +++ b/app/utils/s3Upload.utill.js @@ -1,89 +1,13 @@ -/** - * 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 storage = require("../services/storage/storage.service"); -// app/utils/s3Upload.util.js - -const { PutObjectCommand, GetObjectCommand } = require("@aws-sdk/client-s3"); - -const { getSignedUrl } = require("@aws-sdk/s3-request-presigner"); -const s3 = require("../config/s3.config"); -const { v4: uuidv4 } = require("uuid"); -const path = require("path"); -const redis = require("../config/redisClient"); -const { log } = require("./consoleLog.utill"); - -// Upload + return key (BEST PRACTICE) const uploadToS3 = async (file, usedFor) => { - try { - const fileExtension = path.extname(file.originalname); - const fileName = `${uuidv4()}${fileExtension}`; - let key = `uploads/${fileName}`; - - if (usedFor === "asset") { - key = `assets/${fileName}`; - } - - log("Uploading file to S3 with key:", key); - - const params = { - Bucket: process.env.AWS_S3_BUCKET_NAME, - Key: key, - Body: file.buffer, - ContentType: file.mimetype, - }; - - await s3.send(new PutObjectCommand(params)); - - // Return KEY (not URL) - return key; - } catch (error) { - log("S3 Upload Error:", error); - throw new Error("File upload failed"); - } -}; - -// Generate signed URL with Redis caching -const getSignedFileUrl = async (key) => { - try { - const cacheKey = `s3:url:${key}`; - - // 1. Check cache - const cachedUrl = await redis.get(cacheKey); - - if (cachedUrl) { - log("⚡ Cache HIT - returning signed URL from Redis"); - return cachedUrl; - } - - log("🐢 Cache MISS - generating new signed URL"); - - // 2. Generate signed URL - const command = new GetObjectCommand({ - Bucket: process.env.AWS_S3_BUCKET_NAME, - Key: key, - }); - - const signedUrl = await getSignedUrl(s3, command, { - expiresIn: 3600, // 1 hour - }); - - // 3. Store in Redis (TTL: 55 minutes) - await redis.set(cacheKey, signedUrl, "EX", 3300); - - return signedUrl; - } catch (error) { - log("Signed URL Error:", error); - throw new Error("Failed to generate file URL"); - } + const key = storage.createObjectKey({ ownerId: file.ownerId, mimeType: file.mimetype, purpose: usedFor }); + await storage.uploadBuffer({ buffer: file.buffer, objectKey: key, mimeType: file.mimetype, checksum: file.checksum }); + return key; }; module.exports = { uploadToS3, - getSignedFileUrl, + getSignedFileUrl: storage.createSignedDownloadUrl, + deleteFromS3: storage.deleteObject, }; diff --git a/app/workers/activity.worker.js b/app/workers/activity.worker.js index 39d052a..9611009 100644 --- a/app/workers/activity.worker.js +++ b/app/workers/activity.worker.js @@ -19,7 +19,8 @@ const createActivityWorker = () => { const worker = new Worker( "activity-queue", async (job) => { - await UserActivity.create(job.data); + try { await UserActivity.create(job.data); } + catch (error) { if (error.name === "SequelizeUniqueConstraintError" && job.data.event_id) return { duplicate: true }; throw error; } }, { connection, diff --git a/app/workers/document.worker.js b/app/workers/document.worker.js index 9c69438..b2509b2 100644 --- a/app/workers/document.worker.js +++ b/app/workers/document.worker.js @@ -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; }; diff --git a/app/workers/email.worker.js b/app/workers/email.worker.js new file mode 100644 index 0000000..6867be7 --- /dev/null +++ b/app/workers/email.worker.js @@ -0,0 +1,17 @@ +const { Worker, UnrecoverableError } = require("bullmq"); +const connection = require("../config/redisClient"); +const db = require("../models"); +const { sendMail } = require("../utils/mail.util"); +const { templates } = require("../services/email/email.service"); +const permanent = (error) => ["EENVELOPE", "EMESSAGE"].includes(error.code) || /^5\d\d$/.test(String(error.responseCode || "")); +module.exports = () => { + const worker = new Worker("email-delivery", async (job) => { + const { deliveryId, templateKey, recipient, variables } = job.data; + const delivery = await db.NotificationDelivery.findByPk(deliveryId); + await delivery?.update({ status: "PROCESSING", attemptCount: job.attemptsMade + 1, lastAttemptAt: new Date() }); + try { const result = await sendMail({ to: recipient, subject: templates[templateKey].subject, templateName: templateKey, templateVars: variables }); await delivery?.update({ status: "SENT", sentAt: new Date(), failureCode: null }); return { messageId: result.messageId }; } + catch (error) { if (permanent(error)) throw new UnrecoverableError(error.code || "PERMANENT_REJECTION"); throw error; } + }, { connection, concurrency: Number(process.env.EMAIL_QUEUE_CONCURRENCY || 5) }); + worker.on("failed", async (job, error) => { if (job && (job.attemptsMade >= (job.opts.attempts || 1) || error instanceof UnrecoverableError)) await db.NotificationDelivery.update({ status: "FAILED", failedAt: new Date(), failureCode: String(error.code || error.name || "DELIVERY_FAILED").slice(0, 100) }, { where: { id: job.data.deliveryId } }).catch(() => undefined); console.error("Email delivery failed", { jobId: job?.id, code: error.code || error.name }); }); + return worker; +}; diff --git a/app/workers/index.js b/app/workers/index.js index 71acf50..77e700a 100644 --- a/app/workers/index.js +++ b/app/workers/index.js @@ -22,7 +22,8 @@ const startWorkers = async () => { const createActivityWorker = require("./activity.worker"); const createLogWorker = require("./log.worker"); const createDocumentWorker = require("./document.worker"); - const workers = [createActivityWorker(), createLogWorker(), await createDocumentWorker()]; + const createEmailWorker = require("./email.worker"); + const workers = [createActivityWorker(), createLogWorker(), await createDocumentWorker(), createEmailWorker()]; console.log("All workers started."); let shuttingDown = false; diff --git a/app/workers/log.worker.js b/app/workers/log.worker.js index 7bc4d4d..29f9676 100644 --- a/app/workers/log.worker.js +++ b/app/workers/log.worker.js @@ -25,7 +25,7 @@ const createLogWorker = () => { "logQueue", async (job) => { - const { message, timestamp } = job.data; + const { message, timestamp, level = "info", service = "zumri-api", requestId, metadata } = job.data; const logDate = new Date(timestamp) .toISOString() @@ -36,7 +36,7 @@ const createLogWorker = () => { `app-${logDate}.log` ); - const logLine = `[${timestamp}] ${message}\n`; + const logLine = `${JSON.stringify({ timestamp, level, service, requestId, message, metadata })}\n`; await fs.promises.appendFile( logFilePath, @@ -70,4 +70,4 @@ const createLogWorker = () => { return worker; }; -module.exports = createLogWorker; \ No newline at end of file +module.exports = createLogWorker; diff --git a/cron/notificationCleaning.cron.js b/cron/notificationCleaning.cron.js index 4630de5..702645e 100644 --- a/cron/notificationCleaning.cron.js +++ b/cron/notificationCleaning.cron.js @@ -1,98 +1,17 @@ -/** - * 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. - */ - -// cron/cleanInactiveNotifications.cron.js - const cron = require("node-cron"); +const { Op } = require("sequelize"); +const { sequelize, notification, userNotification } = require("../app/models"); -const { - sequelize, - notification, - userNotification, -} = require("../app/models"); - -function startCleanInactiveNotificationsCron() { - console.log("Inactive Notification Cleanup Cron Started"); - - const task = cron.schedule( - "0 0 * * *", // Run every day at midnight - async () => { - console.log( - `[CRON] Cleaning inactive notifications... ${new Date().toISOString()}`, - ); - - try { - const inactiveNotifications = await notification.findAll({ - where: { - isActive: false, - }, - attributes: ["notification_id"], - }); - - if (!inactiveNotifications.length) { - console.log("[CRON] No inactive notifications found."); - return; - } - - const notificationIds = inactiveNotifications.map( - (item) => item.notification_id, - ); - - const transaction = await sequelize.transaction(); - - try { - // Delete user-notification relationships first - const deletedUserNotifications = await userNotification.destroy({ - where: { - notification_id: notificationIds, - }, - transaction, - }); - - // Delete inactive notifications - const deletedNotifications = await notification.destroy({ - where: { - notification_id: notificationIds, - isActive: false, - }, - transaction, - }); - - await transaction.commit(); - - console.log( - `[CRON] Cleanup completed. Deleted ${deletedNotifications} notifications and ${deletedUserNotifications} user notification records.`, - ); - } catch (err) { - await transaction.rollback(); - - console.error( - "[CRON] Failed to clean inactive notifications:", - err, - ); - } - } catch (err) { - console.error( - "[CRON] Error while checking inactive notifications:", - err, - ); - } - }, - { - scheduled: false, - }, - ); - - task.start(); - - return task; -} - -module.exports = startCleanInactiveNotificationsCron; - +module.exports = function startCleanInactiveNotificationsCron() { + const task = cron.schedule("0 0 * * *", async () => { + const cutoff = new Date(Date.now() - Number(process.env.NOTIFICATION_RETENTION_DAYS || 90) * 86400000); + for (;;) { + const rows = await notification.findAll({ where: { isActive: false, updatedAt: { [Op.lt]: cutoff } }, attributes: ["notification_id"], limit: 500, order: [["updatedAt", "ASC"]] }); + if (!rows.length) break; + const ids = rows.map((row) => row.notification_id); + await sequelize.transaction(async (transaction) => { await userNotification.destroy({ where: { notification_id: { [Op.in]: ids } }, transaction }); await notification.destroy({ where: { notification_id: { [Op.in]: ids }, isActive: false }, transaction }); }); + if (rows.length < 500) break; + } + }, { scheduled: false }); + task.start(); return task; +}; diff --git a/migrations/20260903020000-phase-2-cross-cutting-services.js b/migrations/20260903020000-phase-2-cross-cutting-services.js new file mode 100644 index 0000000..d78d2d3 --- /dev/null +++ b/migrations/20260903020000-phase-2-cross-cutting-services.js @@ -0,0 +1,43 @@ +"use strict"; + +const add = (q, table, column, definition, transaction) => q.addColumn(table, column, definition, { transaction }); +module.exports = { + async up(q, Sequelize) { + await q.sequelize.transaction(async (transaction) => { + await add(q, "uploads", "storage_provider", { type: Sequelize.STRING, allowNull: false, defaultValue: "S3" }, transaction); + await add(q, "uploads", "safe_name", { type: Sequelize.STRING, allowNull: true }, transaction); + await add(q, "uploads", "checksum", { type: Sequelize.STRING(64), allowNull: true }, transaction); + await add(q, "uploads", "owner_type", { type: Sequelize.STRING, allowNull: false, defaultValue: "USER" }, transaction); + await add(q, "uploads", "owner_id", { type: Sequelize.STRING, allowNull: true }, transaction); + await add(q, "uploads", "visibility", { type: Sequelize.ENUM("PRIVATE", "PUBLIC"), allowNull: false, defaultValue: "PRIVATE" }, transaction); + await add(q, "uploads", "status", { type: Sequelize.ENUM("PENDING", "AVAILABLE", "FAILED", "DELETED"), allowNull: false, defaultValue: "AVAILABLE" }, transaction); + await add(q, "uploads", "deletedAt", { type: Sequelize.DATE, allowNull: true }, transaction); + await q.sequelize.query("UPDATE uploads SET owner_id = uploaded_by WHERE owner_id IS NULL", { transaction }); + await q.changeColumn("uploads", "owner_id", { type: Sequelize.STRING, allowNull: false }, { transaction }); + await q.addIndex("uploads", ["owner_id", "status"], { name: "uploads_owner_status_idx", transaction }); + + await q.addIndex("user_notification", ["user_id", "notification_id"], { unique: true, name: "user_notification_user_notification_unique", transaction }); + await q.addIndex("user_notification", ["user_id", "isRead", "createdAt"], { name: "user_notification_inbox_idx", transaction }); + await q.addIndex("notification", ["isActive", "createdAt"], { name: "notification_retention_idx", transaction }); + await q.createTable("notification_deliveries", { + id: { type: Sequelize.STRING, primaryKey: true }, notificationId: { type: Sequelize.STRING, allowNull: true }, userId: { type: Sequelize.STRING, allowNull: true }, + channel: { type: Sequelize.ENUM("EMAIL", "PUSH"), allowNull: false }, recipient: { type: Sequelize.STRING, allowNull: false }, templateKey: { type: Sequelize.STRING, allowNull: false }, + status: { type: Sequelize.ENUM("QUEUED", "PROCESSING", "SENT", "FAILED"), allowNull: false, defaultValue: "QUEUED" }, attemptCount: { type: Sequelize.INTEGER, allowNull: false, defaultValue: 0 }, + lastAttemptAt: Sequelize.DATE, sentAt: Sequelize.DATE, failedAt: Sequelize.DATE, failureCode: Sequelize.STRING, createdAt: { type: Sequelize.DATE, allowNull: false }, updatedAt: { type: Sequelize.DATE, allowNull: false }, + }, { transaction }); + await q.addIndex("notification_deliveries", ["status", "createdAt"], { name: "notification_delivery_status_idx", transaction }); + + for (const [name, type, nullable] of [["event_id", Sequelize.STRING, true],["target_type",Sequelize.STRING,true],["target_id",Sequelize.STRING,true],["ip_address",Sequelize.STRING,true],["user_agent",Sequelize.STRING,true],["request_id",Sequelize.STRING,true],["metadata",Sequelize.JSON,true],["occurred_at",Sequelize.DATE,true]]) await add(q, "UserActivity", name, { type, allowNull: nullable }, transaction); + await q.addIndex("UserActivity", ["event_id"], { unique: true, name: "activity_event_unique", transaction }); + await q.addIndex("UserActivity", ["user_id", "occurred_at"], { name: "activity_actor_time_idx", transaction }); + await q.changeColumn("UserActivity", "user_id", { type: Sequelize.STRING, allowNull: true }, { transaction }); + + await add(q, "Document", "created_by", { type: Sequelize.STRING, allowNull: true }, transaction); await add(q, "Document", "owner_type", { type: Sequelize.STRING, allowNull: false, defaultValue: "USER" }, transaction); await add(q, "Document", "owner_id", { type: Sequelize.STRING, allowNull: true }, transaction); + await add(q, "Document", "storage_upload_id", { type: Sequelize.INTEGER, allowNull: true }, transaction); await add(q, "Document", "job_id", { type: Sequelize.STRING, allowNull: true }, transaction); + await add(q, "Document", "generated_at", { type: Sequelize.DATE, allowNull: true }, transaction); await add(q, "Document", "failed_at", { type: Sequelize.DATE, allowNull: true }, transaction); await add(q, "Document", "failure_code", { type: Sequelize.STRING, allowNull: true }, transaction); await add(q, "Document", "version", { type: Sequelize.INTEGER, allowNull: false, defaultValue: 1 }, transaction); + await q.addIndex("Document", ["job_id"], { unique: true, name: "document_job_unique", transaction }); await q.addIndex("Document", ["owner_id", "status"], { name: "document_owner_status_idx", transaction }); + await q.addIndex("DocumentType", ["doc_type_name"], { unique: true, name: "document_type_name_unique", transaction }); + }); + }, + async down() { throw new Error("Phase 2 migration is forward-only; restore from backup for rollback"); }, +}; diff --git a/tests/unit/cross-cutting-services.test.js b/tests/unit/cross-cutting-services.test.js new file mode 100644 index 0000000..69591a7 --- /dev/null +++ b/tests/unit/cross-cutting-services.test.js @@ -0,0 +1,20 @@ +describe("Phase 2 cross-cutting policies", () => { + test("email templates escape user-controlled HTML", () => { + const { loadTemplate } = require("../../app/utils/mail.util"); + const html = loadTemplate("otp", { firstName: "", otp: "123456" }); + expect(html).toContain("<script>x</script>"); expect(html).not.toContain(""); + }); + test("security notifications bypass disabled marketing preference", () => { + const { channelAllowed } = require("../../app/services/notification/notification-policy.service"); + expect(channelAllowed({ eventType: "LOGIN_OTP", channel: "EMAIL", profile: { notificationsEnabled: false } })).toBe(true); + expect(channelAllowed({ eventType: "MARKETING", channel: "EMAIL", profile: { notificationsEnabled: false } })).toBe(false); + }); + test("queue policies specify bounded retries and retention", () => { + jest.resetModules(); + jest.doMock("../../app/config/redisClient", () => ({})); + jest.doMock("bullmq", () => ({ Queue: jest.fn(function Queue(name, options) { this.name = name; this.opts = options; }) })); + const emailQueue = require("../../app/queues/email.queue"); + expect(emailQueue.opts.defaultJobOptions.attempts).toBe(5); + expect(emailQueue.opts.defaultJobOptions.removeOnComplete).toBeTruthy(); + }); +}); diff --git a/tests/unit/storage.service.test.js b/tests/unit/storage.service.test.js new file mode 100644 index 0000000..6132828 --- /dev/null +++ b/tests/unit/storage.service.test.js @@ -0,0 +1,21 @@ +describe("Phase 2 file validation", () => { + beforeEach(() => jest.resetModules()); + test("detects a valid PNG and creates checksum/safe name", () => { + const { validateUpload } = require("../../app/services/storage/file-validation.service"); + const buffer = Buffer.concat([Buffer.from([0x89,0x50,0x4e,0x47,0x0d,0x0a,0x1a,0x0a]), Buffer.from("payload")]); + const result = validateUpload({ buffer, mimetype: "image/png", originalname: "../unsafe name.svg" }); + expect(result.mimeType).toBe("image/png"); expect(result.safeName).toBe("unsafe_name.png"); expect(result.checksum).toMatch(/^[a-f0-9]{64}$/); + }); + test("rejects claimed MIME that disagrees with magic bytes", () => { + const { validateUpload } = require("../../app/services/storage/file-validation.service"); + expect(() => validateUpload({ buffer: Buffer.from("not an image"), mimetype: "image/png", originalname: "x.png" })).toThrow("does not match"); + }); + test("rejects empty files", () => { + const { validateUpload } = require("../../app/services/storage/file-validation.service"); + expect(() => validateUpload({ buffer: Buffer.alloc(0), mimetype: "image/png", originalname: "x.png" })).toThrow("Empty"); + }); + test("builds controlled owner/date keys without original filenames", () => { + const { createObjectKey } = require("../../app/services/storage/storage.service"); + expect(createObjectKey({ ownerId: "usr_1", mimeType: "application/pdf", purpose: "document", now: new Date("2026-09-03T00:00:00Z"), id: "fixed" })).toBe("documents/usr_1/2026/09/fixed.pdf"); + }); +});