Commit d6c2e6c7 authored by ThinhNC's avatar ThinhNC

Merge branch 'feat/monthly-quota-role-limits-and-cron' into 'develop'

feat(quota): implement monthly quota, role-based limits, reset mechanism and cron scheduler

See merge request !17
parents daf626ee 3d712f1b
...@@ -48,9 +48,12 @@ WORKER_CONCURRENCY=3 ...@@ -48,9 +48,12 @@ WORKER_CONCURRENCY=3
WORKER_JOB_TIMEOUT_MS=300000 WORKER_JOB_TIMEOUT_MS=300000
WORKER_MAX_STALLED_COUNT=1 WORKER_MAX_STALLED_COUNT=1
# Default User Quota Limits
USER_MAX_PAGES=100 USER_MAX_PAGES=100
USER_MAX_JOBS_PER_DAY=10 USER_MAX_JOBS_PER_DAY=10
USER_MAX_CONCURRENT_JOBS=3 USER_MAX_CONCURRENT_JOBS=3
USER_MAX_PAGES_PER_MONTH=1000
USER_MAX_JOBS_PER_MONTH=100
# SMTP Configuration # SMTP Configuration
SMTP_HOST=smtp.gmail.com SMTP_HOST=smtp.gmail.com
......
...@@ -19,6 +19,7 @@ ...@@ -19,6 +19,7 @@
"db:migrate:status": "node scripts/prisma-run.js migrate status", "db:migrate:status": "node scripts/prisma-run.js migrate status",
"db:seed": "node scripts/prisma-run.js db seed -- --tsx prisma/seed.ts", "db:seed": "node scripts/prisma-run.js db seed -- --tsx prisma/seed.ts",
"config:sync": "tsx scripts/sync-system-configs.ts", "config:sync": "tsx scripts/sync-system-configs.ts",
"permission:sync": "tsx scripts/sync-system-permissions.ts",
"lint": "eslint .", "lint": "eslint .",
"format": "prettier --write .", "format": "prettier --write .",
"test": "jest" "test": "jest"
......
-- AlterTable
ALTER TABLE "roles" ADD COLUMN "max_concurrent_jobs_limit" INTEGER NOT NULL DEFAULT 3,
ADD COLUMN "max_jobs_per_day_limit" INTEGER NOT NULL DEFAULT 10,
ADD COLUMN "max_jobs_per_month_limit" INTEGER DEFAULT 100,
ADD COLUMN "max_pages_limit" INTEGER NOT NULL DEFAULT 100,
ADD COLUMN "max_pages_per_month_limit" INTEGER DEFAULT 1000;
-- AlterTable
ALTER TABLE "users" ADD COLUMN "max_jobs_per_month_limit" INTEGER DEFAULT 100,
ADD COLUMN "max_pages_per_month_limit" INTEGER DEFAULT 1000,
ADD COLUMN "quota_reset_at" TIMESTAMP(3);
...@@ -92,6 +92,9 @@ model User { ...@@ -92,6 +92,9 @@ model User {
maxPagesLimit Int @default(100) @map("max_pages_limit") maxPagesLimit Int @default(100) @map("max_pages_limit")
maxJobsPerDayLimit Int @default(10) @map("max_jobs_per_day_limit") maxJobsPerDayLimit Int @default(10) @map("max_jobs_per_day_limit")
maxConcurrentJobsLimit Int @default(3) @map("max_concurrent_jobs_limit") maxConcurrentJobsLimit Int @default(3) @map("max_concurrent_jobs_limit")
maxPagesPerMonthLimit Int? @default(1000) @map("max_pages_per_month_limit")
maxJobsPerMonthLimit Int? @default(100) @map("max_jobs_per_month_limit")
quotaResetAt DateTime? @map("quota_reset_at")
crawlJobs CrawlJob[] crawlJobs CrawlJob[]
crawlSchedules CrawlSchedule[] crawlSchedules CrawlSchedule[]
...@@ -408,6 +411,13 @@ model Role { ...@@ -408,6 +411,13 @@ model Role {
description String? description String?
isSystem Boolean @default(false) @map("is_system") isSystem Boolean @default(false) @map("is_system")
isActive Boolean @default(true) @map("is_active") isActive Boolean @default(true) @map("is_active")
maxPagesLimit Int @default(100) @map("max_pages_limit")
maxJobsPerDayLimit Int @default(10) @map("max_jobs_per_day_limit")
maxConcurrentJobsLimit Int @default(3) @map("max_concurrent_jobs_limit")
maxPagesPerMonthLimit Int? @default(1000) @map("max_pages_per_month_limit")
maxJobsPerMonthLimit Int? @default(100) @map("max_jobs_per_month_limit")
createdAt DateTime @default(now()) @map("created_at") createdAt DateTime @default(now()) @map("created_at")
updatedAt DateTime @updatedAt @map("updated_at") updatedAt DateTime @updatedAt @map("updated_at")
......
import "dotenv/config";
import { PermissionRepository } from "../src/modules/permissions/permission.repository";
import { authorizationCache } from "../src/common/helpers/authorization-cache.helper";
import { prisma } from "../src/database/prisma.client";
async function main() {
console.log("Synchronizing system permissions and roles with PostgreSQL database...");
const repository = new PermissionRepository();
await repository.ensureSystemPermissions();
authorizationCache.invalidateAll();
const totalPerms = await prisma.permission.count();
const totalRoles = await prisma.role.count();
const totalRolePerms = await prisma.rolePermission.count();
console.log(`\n============================ SUMMARY ============================`);
console.log(`[✔] Total permissions in DB: ${totalPerms}`);
console.log(`[✔] Total roles in DB: ${totalRoles}`);
console.log(`[✔] Total role-permission bindings: ${totalRolePerms}`);
console.log("System permissions synchronized successfully!");
}
main()
.catch((err) => {
console.error("Failed to sync system permissions:", err);
process.exit(1);
})
.finally(async () => {
await prisma.$disconnect();
});
...@@ -36,6 +36,11 @@ export const AUDIT_ACTIONS = { ...@@ -36,6 +36,11 @@ export const AUDIT_ACTIONS = {
SYSTEM_CONFIG_UPDATED: "SYSTEM_CONFIG_UPDATED", SYSTEM_CONFIG_UPDATED: "SYSTEM_CONFIG_UPDATED",
SYSTEM_CONFIG_TOGGLED: "SYSTEM_CONFIG_TOGGLED", SYSTEM_CONFIG_TOGGLED: "SYSTEM_CONFIG_TOGGLED",
SYSTEM_CONFIG_DELETED: "SYSTEM_CONFIG_DELETED", SYSTEM_CONFIG_DELETED: "SYSTEM_CONFIG_DELETED",
CRON_JOB_TOGGLED: "CRON_JOB_TOGGLED",
CRON_JOB_TRIGGERED: "CRON_JOB_TRIGGERED",
CRON_JOB_EXECUTED: "CRON_JOB_EXECUTED",
USER_QUOTA_RESET: "USER_QUOTA_RESET",
ROLE_QUOTA_RESET: "ROLE_QUOTA_RESET",
} as const; } as const;
export type AuditAction = (typeof AUDIT_ACTIONS)[keyof typeof AUDIT_ACTIONS]; export type AuditAction = (typeof AUDIT_ACTIONS)[keyof typeof AUDIT_ACTIONS];
export const CRON_JOB_NAMES = {
CLEANUP_AUDIT_LOGS: "cleanup-audit-logs",
CLEANUP_UNCONFIRMED_UPLOADS: "cleanup-unconfirmed-uploads",
CLEANUP_EXPIRED_TOKENS: "cleanup-expired-tokens",
DAILY_SUMMARY_DIGEST: "daily-summary-digest",
WEEKLY_SUMMARY_DIGEST: "weekly-summary-digest",
} as const;
export type CronJobName = (typeof CRON_JOB_NAMES)[keyof typeof CRON_JOB_NAMES];
export const CRON_JOB_STATUS = {
READY: "READY",
RUNNING: "RUNNING",
SUCCESS: "SUCCESS",
FAILED: "FAILED",
} as const;
export type CronJobStatus =
(typeof CRON_JOB_STATUS)[keyof typeof CRON_JOB_STATUS];
export const DEFAULT_CRON_TIMEZONE = "Asia/Ho_Chi_Minh" as const;
export const CRON_QUEUE_NAME = "cron-scheduler-queue" as const;
export const CRON_SYSTEM_CONFIG_KEY = "CRON_JOB_STATUSES" as const;
export const DEFAULT_AUDIT_LOG_RETENTION_DAYS = 30;
export const DEFAULT_UNCONFIRMED_UPLOAD_MAX_AGE_HOURS = 24;
export interface CronScheduleConfig {
cron: string;
description: string;
defaultParams?: Record<string, unknown>;
}
export const DEFAULT_CRON_SCHEDULES: Record<CronJobName, CronScheduleConfig> = {
[CRON_JOB_NAMES.CLEANUP_AUDIT_LOGS]: {
cron: "0 2 * * *",
description: "Dọn dẹp các bản ghi nhật ký kiểm toán cũ hơn số ngày quy định",
defaultParams: {
retentionDays: DEFAULT_AUDIT_LOG_RETENTION_DAYS,
},
},
[CRON_JOB_NAMES.CLEANUP_UNCONFIRMED_UPLOADS]: {
cron: "0 3 * * *",
description: "Quét và dọn dẹp các tệp tin tải lên mồ côi hoặc xuất file tạm quá hạn",
defaultParams: {
maxAgeHours: DEFAULT_UNCONFIRMED_UPLOAD_MAX_AGE_HOURS,
},
},
[CRON_JOB_NAMES.CLEANUP_EXPIRED_TOKENS]: {
cron: "0 4 * * *",
description: "Dọn dẹp các mã phiên xác thực và token đăng nhập đã hết hạn",
},
[CRON_JOB_NAMES.DAILY_SUMMARY_DIGEST]: {
cron: "0 8 * * *",
description: "Tổng hợp chỉ số hoạt động ngày hôm trước và gửi email báo cáo cho quản trị viên",
},
[CRON_JOB_NAMES.WEEKLY_SUMMARY_DIGEST]: {
cron: "0 8 * * 1",
description: "Tổng hợp chỉ số hiệu suất hệ thống 7 ngày gần nhất và gửi email báo cáo tuần",
},
};
...@@ -13,3 +13,5 @@ export * from "./webhook.constant"; ...@@ -13,3 +13,5 @@ export * from "./webhook.constant";
export * from "./system-role.constant"; export * from "./system-role.constant";
export * from "./permission.constant"; export * from "./permission.constant";
export * from "./system-config.constant"; export * from "./system-config.constant";
export * from "./cron.constant";
...@@ -72,6 +72,10 @@ export const PERMISSIONS = { ...@@ -72,6 +72,10 @@ export const PERMISSIONS = {
// System Configs // System Configs
SYSTEM_CONFIG_READ: "system_configs.read", SYSTEM_CONFIG_READ: "system_configs.read",
SYSTEM_CONFIG_MANAGE: "system_configs.manage", SYSTEM_CONFIG_MANAGE: "system_configs.manage",
// Cron Jobs
CRON_JOB_READ: "cron_jobs.read",
CRON_JOB_MANAGE: "cron_jobs.manage",
} as const; } as const;
export type PermissionSlug = (typeof PERMISSIONS)[keyof typeof PERMISSIONS]; export type PermissionSlug = (typeof PERMISSIONS)[keyof typeof PERMISSIONS];
...@@ -493,6 +497,23 @@ export const SYSTEM_PERMISSIONS_CATALOG: PermissionDefinition[] = [ ...@@ -493,6 +497,23 @@ export const SYSTEM_PERMISSIONS_CATALOG: PermissionDefinition[] = [
action: "manage", action: "manage",
isSystem: true, isSystem: true,
}, },
// Cron Jobs
{
name: "View Cron Jobs",
slug: PERMISSIONS.CRON_JOB_READ,
description: "Xem danh sách tác vụ định kỳ và lịch chạy nền",
resource: "cron_jobs",
action: "read",
isSystem: true,
},
{
name: "Manage Cron Jobs",
slug: PERMISSIONS.CRON_JOB_MANAGE,
description: "Kích hoạt chạy thủ công và bật/tắt lịch chạy tự động",
resource: "cron_jobs",
action: "manage",
isSystem: true,
},
]; ];
export const SYSTEM_ROLE_DEFAULT_PERMISSIONS: Record< export const SYSTEM_ROLE_DEFAULT_PERMISSIONS: Record<
...@@ -549,6 +570,8 @@ export const SYSTEM_ROLE_DEFAULT_PERMISSIONS: Record< ...@@ -549,6 +570,8 @@ export const SYSTEM_ROLE_DEFAULT_PERMISSIONS: Record<
PERMISSIONS.DASHBOARD_READ_ALL, PERMISSIONS.DASHBOARD_READ_ALL,
PERMISSIONS.SYSTEM_CONFIG_READ, PERMISSIONS.SYSTEM_CONFIG_READ,
PERMISSIONS.SYSTEM_CONFIG_MANAGE, PERMISSIONS.SYSTEM_CONFIG_MANAGE,
PERMISSIONS.CRON_JOB_READ,
PERMISSIONS.CRON_JOB_MANAGE,
], ],
[SYSTEM_ROLE_SLUGS.CRAWLER_USER]: [ [SYSTEM_ROLE_SLUGS.CRAWLER_USER]: [
PERMISSIONS.CRAWL_JOBS_CREATE, PERMISSIONS.CRAWL_JOBS_CREATE,
......
...@@ -295,6 +295,20 @@ export const DEFAULT_SYSTEM_CONFIGS: readonly DefaultSystemConfigItem[] = [ ...@@ -295,6 +295,20 @@ export const DEFAULT_SYSTEM_CONFIGS: readonly DefaultSystemConfigItem[] = [
isPublic: false, isPublic: false,
description: "Hạn ngạch số tác vụ cào được phép chạy song song cho mỗi tài khoản (USER_MAX_CONCURRENT_JOBS)", description: "Hạn ngạch số tác vụ cào được phép chạy song song cho mỗi tài khoản (USER_MAX_CONCURRENT_JOBS)",
}, },
{
key: "quota.user_max_pages_per_month",
value: 1000,
category: SYSTEM_CONFIG_CATEGORY.SECURITY,
isPublic: false,
description: "Hạn ngạch số trang cào tối đa trong một tháng cho mỗi tài khoản (USER_MAX_PAGES_PER_MONTH)",
},
{
key: "quota.user_max_jobs_per_month",
value: 100,
category: SYSTEM_CONFIG_CATEGORY.SECURITY,
isPublic: false,
description: "Hạn ngạch số tác vụ cào tối đa trong một tháng cho mỗi tài khoản (USER_MAX_JOBS_PER_MONTH)",
},
{ {
key: "rate_limit.window_ms", key: "rate_limit.window_ms",
value: 900000, value: 900000,
......
...@@ -18,6 +18,8 @@ export const ERROR_CODE = { ...@@ -18,6 +18,8 @@ export const ERROR_CODE = {
QUOTA_MAX_PAGES_EXCEEDED: "QUOTA_MAX_PAGES_EXCEEDED", QUOTA_MAX_PAGES_EXCEEDED: "QUOTA_MAX_PAGES_EXCEEDED",
QUOTA_JOBS_PER_DAY_EXCEEDED: "QUOTA_JOBS_PER_DAY_EXCEEDED", QUOTA_JOBS_PER_DAY_EXCEEDED: "QUOTA_JOBS_PER_DAY_EXCEEDED",
QUOTA_CONCURRENT_JOBS_EXCEEDED: "QUOTA_CONCURRENT_JOBS_EXCEEDED", QUOTA_CONCURRENT_JOBS_EXCEEDED: "QUOTA_CONCURRENT_JOBS_EXCEEDED",
QUOTA_JOBS_PER_MONTH_EXCEEDED: "QUOTA_JOBS_PER_MONTH_EXCEEDED",
QUOTA_MAX_PAGES_PER_MONTH_EXCEEDED: "QUOTA_MAX_PAGES_PER_MONTH_EXCEEDED",
EXPORT_NOT_FOUND: "EXPORT_NOT_FOUND", EXPORT_NOT_FOUND: "EXPORT_NOT_FOUND",
EXPORT_FILE_MISSING: "EXPORT_FILE_MISSING", EXPORT_FILE_MISSING: "EXPORT_FILE_MISSING",
UNSUPPORTED_EXPORT_TYPE: "UNSUPPORTED_EXPORT_TYPE", UNSUPPORTED_EXPORT_TYPE: "UNSUPPORTED_EXPORT_TYPE",
......
...@@ -91,6 +91,14 @@ export const envConfig = { ...@@ -91,6 +91,14 @@ export const envConfig = {
process.env.USER_MAX_CONCURRENT_JOBS || "3", process.env.USER_MAX_CONCURRENT_JOBS || "3",
10, 10,
), ),
defaultMaxPagesPerMonth: parseInt(
process.env.USER_MAX_PAGES_PER_MONTH || "1000",
10,
),
defaultMaxJobsPerMonth: parseInt(
process.env.USER_MAX_JOBS_PER_MONTH || "100",
10,
),
}, },
cors: { cors: {
allowedOrigins: ( allowedOrigins: (
......
...@@ -321,6 +321,8 @@ export const swaggerPaths: Record<string, any> = { ...@@ -321,6 +321,8 @@ export const swaggerPaths: Record<string, any> = {
example: 3, example: 3,
}, },
totalPagesCrawled: { type: "integer", example: 45 }, totalPagesCrawled: { type: "integer", example: 45 },
pagesCrawledToday: { type: "integer", example: 0 },
pagesRemainingToday: { type: "integer", example: 100 },
}, },
}, },
resetAt: { resetAt: {
......
...@@ -472,6 +472,14 @@ ...@@ -472,6 +472,14 @@
"totalPagesCrawled": { "totalPagesCrawled": {
"type": "integer", "type": "integer",
"example": 45 "example": 45
},
"pagesCrawledToday": {
"type": "integer",
"example": 0
},
"pagesRemainingToday": {
"type": "integer",
"example": 100
} }
} }
}, },
...@@ -1286,6 +1294,26 @@ ...@@ -1286,6 +1294,26 @@
"summary": "Xóa người dùng" "summary": "Xóa người dùng"
} }
}, },
"/users/{id}/reset-quota": {
"post": {
"description": "",
"parameters": [
{
"name": "id",
"in": "path",
"required": true,
"schema": {
"type": "string"
}
}
],
"responses": {
"default": {
"description": ""
}
}
}
},
"/users/{id}/roles": { "/users/{id}/roles": {
"get": { "get": {
"description": "", "description": "",
...@@ -1656,6 +1684,26 @@ ...@@ -1656,6 +1684,26 @@
"summary": "Xóa Role" "summary": "Xóa Role"
} }
}, },
"/roles/{id}/reset-quota": {
"post": {
"description": "",
"parameters": [
{
"name": "id",
"in": "path",
"required": true,
"schema": {
"type": "string"
}
}
],
"responses": {
"default": {
"description": ""
}
}
}
},
"/roles/{id}/permissions": { "/roles/{id}/permissions": {
"get": { "get": {
"description": "", "description": "",
...@@ -2183,6 +2231,93 @@ ...@@ -2183,6 +2231,93 @@
"summary": "Bật/tắt nhanh Feature Flag" "summary": "Bật/tắt nhanh Feature Flag"
} }
}, },
"/cron/jobs": {
"get": {
"description": "",
"parameters": [
{
"name": "search",
"in": "query",
"schema": {
"type": "string"
}
}
],
"responses": {
"200": {
"description": "OK"
}
}
}
},
"/cron/jobs/{jobName}/trigger": {
"post": {
"description": "",
"parameters": [
{
"name": "jobName",
"in": "path",
"required": true,
"schema": {
"type": "string"
}
}
],
"responses": {
"200": {
"description": "OK"
}
},
"requestBody": {
"content": {
"application/json": {
"schema": {
"type": "object",
"properties": {
"params": {
"example": "any"
}
}
}
}
}
}
}
},
"/cron/jobs/{jobName}/toggle": {
"patch": {
"description": "",
"parameters": [
{
"name": "jobName",
"in": "path",
"required": true,
"schema": {
"type": "string"
}
}
],
"responses": {
"200": {
"description": "OK"
}
},
"requestBody": {
"content": {
"application/json": {
"schema": {
"type": "object",
"properties": {
"enabled": {
"example": "any"
}
}
}
}
}
}
}
},
"/dashboard/stats": { "/dashboard/stats": {
"get": { "get": {
"description": "Thống kê tổng hợp số lượng crawl jobs theo trạng thái, số trang đã crawl, số lịch crawl đang chạy và tổng số tệp export.", "description": "Thống kê tổng hợp số lượng crawl jobs theo trạng thái, số trang đã crawl, số lịch crawl đang chạy và tổng số tệp export.",
......
...@@ -115,4 +115,36 @@ describe("Permission Middleware", () => { ...@@ -115,4 +115,36 @@ describe("Permission Middleware", () => {
); );
}); });
}); });
describe("Role-based permission inheritance and Super Admin bypass", () => {
it("should grant full access to SUPER_ADMIN even if permission array is empty", async () => {
mockReq.user = {
id: "super-1",
email: "super@example.com",
role: "SUPER_ADMIN",
roles: ["super_admin"],
permissions: [],
} as any;
const middleware = requirePermission(PERMISSIONS.CRON_JOB_READ);
await middleware(mockReq as Request, mockRes as Response, mockNext);
expect(mockNext).toHaveBeenCalledWith();
});
it("should grant CRON_JOB_READ to ADMIN role even if permissions array was missing it", async () => {
mockReq.user = {
id: "admin-1",
email: "admin@example.com",
role: "ADMIN",
roles: ["admin"],
permissions: [PERMISSIONS.USERS_READ], // Missing CRON_JOB_READ in DB array
} as any;
const middleware = requirePermission(PERMISSIONS.CRON_JOB_READ);
await middleware(mockReq as Request, mockRes as Response, mockNext);
expect(mockNext).toHaveBeenCalledWith();
});
});
}); });
import { Request, Response, NextFunction } from "express"; import { Request, Response, NextFunction } from "express";
import { import {
PERMISSIONS,
PermissionSlug, PermissionSlug,
SYSTEM_ROLE_DEFAULT_PERMISSIONS, SYSTEM_ROLE_DEFAULT_PERMISSIONS,
} from "../common/constants/permission.constant"; } from "../common/constants/permission.constant";
import { SystemRoleSlug } from "../common/constants/system-role.constant"; import {
SYSTEM_ROLE_SLUGS,
SystemRoleSlug,
} from "../common/constants/system-role.constant";
import { AppError } from "../common/errors/app-error"; import { AppError } from "../common/errors/app-error";
import { ERROR_CODE } from "../common/errors/error-code"; import { ERROR_CODE } from "../common/errors/error-code";
import { PermissionService } from "../modules/permissions/permission.service"; import { permissionService } from "../modules/permissions/permission.service";
const permissionService = new PermissionService();
async function resolveUserPermissions(req: Request): Promise<string[]> { async function resolveUserPermissions(req: Request): Promise<string[]> {
if (req.user.permissions && Array.isArray(req.user.permissions)) { let permissions = req.user.permissions;
return req.user.permissions; if (!permissions || !Array.isArray(permissions)) {
permissions = await permissionService.getUserPermissions(req.user.id);
} }
const permissions = await permissionService.getUserPermissions(req.user.id);
if (!req.user.roles) { if (!req.user.roles) {
req.user.roles = await permissionService.getUserRoles(req.user.id); req.user.roles = await permissionService.getUserRoles(req.user.id);
} }
// Super Admin has all system permissions
const isSuperAdmin =
req.user.roles?.includes(SYSTEM_ROLE_SLUGS.SUPER_ADMIN) ||
String(req.user.role).toLowerCase() === "super_admin";
if (isSuperAdmin) {
const allPerms = Object.values(PERMISSIONS) as string[];
req.user.permissions = allPerms;
return allPerms;
}
// Admin automatically inherits all default admin permissions merged with any assigned permissions
const isAdmin =
req.user.role === "ADMIN" ||
String(req.user.role).toLowerCase() === "admin" ||
req.user.roles?.includes(SYSTEM_ROLE_SLUGS.ADMIN);
if (isAdmin) {
const adminDefaults =
SYSTEM_ROLE_DEFAULT_PERMISSIONS[SYSTEM_ROLE_SLUGS.ADMIN] || [];
const merged = Array.from(new Set([...permissions, ...adminDefaults]));
req.user.permissions = merged;
return merged;
}
if (permissions.length === 0 && req.user.role) { if (permissions.length === 0 && req.user.role) {
const normalizedRole = String(req.user.role).toLowerCase() as SystemRoleSlug;
const defaultPerms = const defaultPerms =
SYSTEM_ROLE_DEFAULT_PERMISSIONS[req.user.role as SystemRoleSlug] || []; SYSTEM_ROLE_DEFAULT_PERMISSIONS[normalizedRole] ||
SYSTEM_ROLE_DEFAULT_PERMISSIONS[req.user.role as SystemRoleSlug] ||
[];
req.user.permissions = defaultPerms; req.user.permissions = defaultPerms;
return defaultPerms; return defaultPerms;
} }
......
...@@ -14,6 +14,9 @@ describe("AuthService getUsage and avatarUrl", () => { ...@@ -14,6 +14,9 @@ describe("AuthService getUsage and avatarUrl", () => {
maxPagesLimit: 100, maxPagesLimit: 100,
maxJobsPerDayLimit: 10, maxJobsPerDayLimit: 10,
maxConcurrentJobsLimit: 3, maxConcurrentJobsLimit: 3,
maxPagesPerMonthLimit: 1000,
maxJobsPerMonthLimit: 100,
quotaResetAt: null,
createdAt: new Date(), createdAt: new Date(),
}; };
...@@ -24,15 +27,16 @@ describe("AuthService getUsage and avatarUrl", () => { ...@@ -24,15 +27,16 @@ describe("AuthService getUsage and avatarUrl", () => {
}; };
(service as any).repository = repository; (service as any).repository = repository;
( (CrawlJobRepository.prototype.countJobsSince as jest.Mock)
CrawlJobRepository.prototype.countJobsSince as jest.Mock .mockResolvedValueOnce(4)
).mockResolvedValue(4); .mockResolvedValueOnce(15);
( (
CrawlJobRepository.prototype.countConcurrentJobs as jest.Mock CrawlJobRepository.prototype.countConcurrentJobs as jest.Mock
).mockResolvedValue(1); ).mockResolvedValue(1);
( (CrawlJobRepository.prototype.sumPagesCrawledByUser as jest.Mock)
CrawlJobRepository.prototype.sumPagesCrawledByUser as jest.Mock .mockResolvedValueOnce(125)
).mockResolvedValue(125); .mockResolvedValueOnce(25)
.mockResolvedValueOnce(75);
const usage = await service.getUsage("user-123"); const usage = await service.getUsage("user-123");
...@@ -40,6 +44,8 @@ describe("AuthService getUsage and avatarUrl", () => { ...@@ -40,6 +44,8 @@ describe("AuthService getUsage and avatarUrl", () => {
maxPagesLimit: 100, maxPagesLimit: 100,
maxJobsPerDayLimit: 10, maxJobsPerDayLimit: 10,
maxConcurrentJobsLimit: 3, maxConcurrentJobsLimit: 3,
maxPagesPerMonthLimit: 1000,
maxJobsPerMonthLimit: 100,
}); });
expect(usage.usage).toEqual({ expect(usage.usage).toEqual({
jobsUsedToday: 4, jobsUsedToday: 4,
...@@ -47,8 +53,16 @@ describe("AuthService getUsage and avatarUrl", () => { ...@@ -47,8 +53,16 @@ describe("AuthService getUsage and avatarUrl", () => {
concurrentJobsRunning: 1, concurrentJobsRunning: 1,
concurrentJobsAvailable: 2, concurrentJobsAvailable: 2,
totalPagesCrawled: 125, totalPagesCrawled: 125,
pagesCrawledToday: 25,
pagesRemainingToday: 75,
jobsUsedThisMonth: 15,
jobsRemainingThisMonth: 85,
pagesCrawledThisMonth: 75,
pagesRemainingThisMonth: 925,
}); });
expect(usage.resetAt).toBeDefined(); expect(usage.resetAt).toBeDefined();
expect(usage.monthlyResetAt).toBeDefined();
expect(usage.quotaResetAt).toBeNull();
}); });
it("updates fullName via updateMe", async () => { it("updates fullName via updateMe", async () => {
......
...@@ -49,6 +49,8 @@ export interface UserUsageDto { ...@@ -49,6 +49,8 @@ export interface UserUsageDto {
maxPagesLimit: number; maxPagesLimit: number;
maxJobsPerDayLimit: number; maxJobsPerDayLimit: number;
maxConcurrentJobsLimit: number; maxConcurrentJobsLimit: number;
maxPagesPerMonthLimit?: number | null;
maxJobsPerMonthLimit?: number | null;
}; };
usage: { usage: {
jobsUsedToday: number; jobsUsedToday: number;
...@@ -56,8 +58,16 @@ export interface UserUsageDto { ...@@ -56,8 +58,16 @@ export interface UserUsageDto {
concurrentJobsRunning: number; concurrentJobsRunning: number;
concurrentJobsAvailable: number; concurrentJobsAvailable: number;
totalPagesCrawled: number; totalPagesCrawled: number;
pagesCrawledToday: number;
pagesRemainingToday: number;
jobsUsedThisMonth: number;
jobsRemainingThisMonth?: number | null;
pagesCrawledThisMonth: number;
pagesRemainingThisMonth?: number | null;
}; };
resetAt: string; resetAt: string;
monthlyResetAt?: string;
quotaResetAt?: string | null;
} }
export interface ChangePasswordDto { export interface ChangePasswordDto {
......
...@@ -512,6 +512,31 @@ export class AuthService { ...@@ -512,6 +512,31 @@ export class AuthService {
DEFAULT_TIMEZONE, DEFAULT_TIMEZONE,
); );
const startOfMonth = createUtcDateFromZonedParts(
nowZoned.year,
nowZoned.month,
1,
0,
0,
DEFAULT_TIMEZONE,
);
const nextMonth = createUtcDateFromZonedParts(
nowZoned.month === 12 ? nowZoned.year + 1 : nowZoned.year,
nowZoned.month === 12 ? 1 : nowZoned.month + 1,
1,
0,
0,
DEFAULT_TIMEZONE,
);
const quotaResetAt = (user as any).quotaResetAt
? new Date((user as any).quotaResetAt)
: null;
const effectiveDailySince =
quotaResetAt && quotaResetAt > startOfDay ? quotaResetAt : startOfDay;
const effectiveMonthlySince =
quotaResetAt && quotaResetAt > startOfMonth ? quotaResetAt : startOfMonth;
const twoHoursAgo = new Date(); const twoHoursAgo = new Date();
twoHoursAgo.setHours(twoHoursAgo.getHours() - 2); twoHoursAgo.setHours(twoHoursAgo.getHours() - 2);
const activeStatuses = [ const activeStatuses = [
...@@ -521,19 +546,41 @@ export class AuthService { ...@@ -521,19 +546,41 @@ export class AuthService {
JOB_STATUS.PROCESSING_EXPORT, JOB_STATUS.PROCESSING_EXPORT,
]; ];
const [jobsTodayCount, concurrentJobsCount, totalPages] = await Promise.all( const [
[ jobsTodayCount,
crawlJobRepo.countJobsSince(userId, startOfDay), concurrentJobsCount,
totalPages,
pagesToday,
jobsMonthCount,
pagesMonth,
] = await Promise.all([
crawlJobRepo.countJobsSince(userId, effectiveDailySince),
crawlJobRepo.countConcurrentJobs(userId, activeStatuses, twoHoursAgo), crawlJobRepo.countConcurrentJobs(userId, activeStatuses, twoHoursAgo),
crawlJobRepo.sumPagesCrawledByUser(userId), crawlJobRepo.sumPagesCrawledByUser(userId),
], crawlJobRepo.sumPagesCrawledByUser(userId, effectiveDailySince),
); crawlJobRepo.countJobsSince(userId, effectiveMonthlySince),
crawlJobRepo.sumPagesCrawledByUser(userId, effectiveMonthlySince),
]);
const maxPagesPerMonthLimit = (user as any).maxPagesPerMonthLimit ?? null;
const maxJobsPerMonthLimit = (user as any).maxJobsPerMonthLimit ?? null;
const jobsRemainingThisMonth =
maxJobsPerMonthLimit !== null
? Math.max(0, maxJobsPerMonthLimit - jobsMonthCount)
: null;
const pagesRemainingThisMonth =
maxPagesPerMonthLimit !== null
? Math.max(0, maxPagesPerMonthLimit - pagesMonth)
: null;
return { return {
quota: { quota: {
maxPagesLimit: user.maxPagesLimit, maxPagesLimit: user.maxPagesLimit,
maxJobsPerDayLimit: user.maxJobsPerDayLimit, maxJobsPerDayLimit: user.maxJobsPerDayLimit,
maxConcurrentJobsLimit: user.maxConcurrentJobsLimit, maxConcurrentJobsLimit: user.maxConcurrentJobsLimit,
maxPagesPerMonthLimit,
maxJobsPerMonthLimit,
}, },
usage: { usage: {
jobsUsedToday: jobsTodayCount, jobsUsedToday: jobsTodayCount,
...@@ -547,8 +594,16 @@ export class AuthService { ...@@ -547,8 +594,16 @@ export class AuthService {
user.maxConcurrentJobsLimit - concurrentJobsCount, user.maxConcurrentJobsLimit - concurrentJobsCount,
), ),
totalPagesCrawled: totalPages, totalPagesCrawled: totalPages,
pagesCrawledToday: pagesToday,
pagesRemainingToday: Math.max(0, user.maxPagesLimit - pagesToday),
jobsUsedThisMonth: jobsMonthCount,
jobsRemainingThisMonth,
pagesCrawledThisMonth: pagesMonth,
pagesRemainingThisMonth,
}, },
resetAt: nextDay.toISOString(), resetAt: nextDay.toISOString(),
monthlyResetAt: nextMonth.toISOString(),
quotaResetAt: quotaResetAt ? quotaResetAt.toISOString() : null,
}; };
} }
......
...@@ -380,9 +380,12 @@ export class CrawlJobRepository { ...@@ -380,9 +380,12 @@ export class CrawlJobRepository {
}); });
} }
async sumPagesCrawledByUser(userId: string): Promise<number> { async sumPagesCrawledByUser(userId: string, since?: Date): Promise<number> {
const aggregate = await prisma.crawlJob.aggregate({ const aggregate = await prisma.crawlJob.aggregate({
where: { userId }, where: {
userId,
...(since ? { createdAt: { gte: since } } : {}),
},
_sum: { totalPages: true }, _sum: { totalPages: true },
}); });
const totalPages = aggregate._sum.totalPages ?? 0; const totalPages = aggregate._sum.totalPages ?? 0;
...@@ -394,6 +397,7 @@ export class CrawlJobRepository { ...@@ -394,6 +397,7 @@ export class CrawlJobRepository {
userId, userId,
status: JOB_STATUS.FAILED, status: JOB_STATUS.FAILED,
successPages: 0, successPages: 0,
...(since ? { createdAt: { gte: since } } : {}),
}, },
select: { select: {
totalPages: true, totalPages: true,
......
import { CronService } from "../cron.service";
import { CronRepository } from "../cron.repository";
import { IStorageService } from "../../../common/storage/storage.interface";
import { MailService } from "../../mail/mail.service";
import { CronQueueService } from "../../../queues/cron.queue";
import {
CRON_JOB_NAMES,
CRON_JOB_STATUS,
DEFAULT_CRON_SCHEDULES,
} from "../../../common/constants/cron.constant";
import { AUDIT_ACTIONS } from "../../../common/constants/audit-action.constant";
describe("CronService", () => {
let cronService: CronService;
let mockRepository: jest.Mocked<CronRepository>;
let mockStorageService: jest.Mocked<IStorageService>;
let mockMailService: jest.Mocked<MailService>;
let mockQueueService: jest.Mocked<CronQueueService>;
beforeEach(() => {
mockRepository = {
getJobStatuses: jest.fn().mockResolvedValue({
[CRON_JOB_NAMES.CLEANUP_AUDIT_LOGS]: true,
[CRON_JOB_NAMES.DAILY_SUMMARY_DIGEST]: false,
}),
setJobStatus: jest.fn().mockResolvedValue({}),
getLatestRunsForJobs: jest.fn().mockResolvedValue(
new Map([
[
CRON_JOB_NAMES.CLEANUP_AUDIT_LOGS,
{
lastRun: new Date("2026-09-08T02:00:00.000Z"),
status: CRON_JOB_STATUS.SUCCESS,
durationMs: 120,
},
],
]),
),
createAuditLog: jest.fn().mockResolvedValue({} as any),
cleanupAuditLogs: jest.fn().mockResolvedValue(45),
cleanupExpiredTokens: jest.fn().mockResolvedValue({ deletedRefreshTokens: 12 }),
cleanupExpiredExports: jest.fn().mockResolvedValue({
deletedCount: 2,
filePaths: ["exports/test1.zip", "exports/test2.zip"],
}),
getActiveAvatarUrls: jest.fn().mockResolvedValue([]),
getDigestStats: jest.fn().mockResolvedValue({
newUsers: 5,
crawlJobsTotal: 10,
crawlJobsCompleted: 8,
crawlJobsFailed: 2,
crawledPages: 150,
exportsGenerated: 4,
webhookDeliveries: 10,
auditLogsRecorded: 80,
}),
getAdminEmails: jest.fn().mockResolvedValue(["admin@datacrawler.com"]),
} as unknown as jest.Mocked<CronRepository>;
mockStorageService = {
uploadFile: jest.fn(),
uploadStream: jest.fn(),
downloadFile: jest.fn(),
getReadStream: jest.fn(),
deleteFile: jest.fn().mockResolvedValue(undefined),
exists: jest.fn().mockResolvedValue(true),
};
mockMailService = {
sendDigestEmail: jest.fn().mockResolvedValue(undefined),
sendPasswordResetEmail: jest.fn(),
sendVerificationEmail: jest.fn(),
sendDeactivationEmail: jest.fn(),
} as unknown as jest.Mocked<MailService>;
mockQueueService = {
getQueue: jest.fn(),
registerDefaultSchedulers: jest.fn().mockResolvedValue(undefined),
enableJobScheduler: jest.fn().mockResolvedValue(undefined),
disableJobScheduler: jest.fn().mockResolvedValue(undefined),
} as unknown as jest.Mocked<CronQueueService>;
cronService = new CronService(
mockRepository,
mockStorageService,
mockMailService,
mockQueueService,
);
});
describe("listJobs", () => {
it("should return all configured jobs with correct status and history", async () => {
const jobs = await cronService.listJobs();
expect(jobs.length).toBe(Object.keys(DEFAULT_CRON_SCHEDULES).length);
const auditLogJob = jobs.find(
(j) => j.name === CRON_JOB_NAMES.CLEANUP_AUDIT_LOGS,
);
expect(auditLogJob).toBeDefined();
expect(auditLogJob?.isEnabled).toBe(true);
expect(auditLogJob?.lastStatus).toBe(CRON_JOB_STATUS.SUCCESS);
expect(auditLogJob?.lastDurationMs).toBe(120);
const digestJob = jobs.find(
(j) => j.name === CRON_JOB_NAMES.DAILY_SUMMARY_DIGEST,
);
expect(digestJob?.isEnabled).toBe(false);
});
it("should filter jobs when search parameter is provided", async () => {
const jobs = await cronService.listJobs("audit");
expect(jobs.length).toBe(1);
expect(jobs[0].name).toBe(CRON_JOB_NAMES.CLEANUP_AUDIT_LOGS);
});
});
describe("toggleJob", () => {
it("should update job status in DB, sync BullMQ, and record audit log", async () => {
const result = await cronService.toggleJob(
CRON_JOB_NAMES.CLEANUP_AUDIT_LOGS,
false,
{ actorId: "admin-1", ipAddress: "127.0.0.1", userAgent: "Jest" },
);
expect(mockRepository.setJobStatus).toHaveBeenCalledWith(
CRON_JOB_NAMES.CLEANUP_AUDIT_LOGS,
false,
);
expect(mockQueueService.disableJobScheduler).toHaveBeenCalledWith(
CRON_JOB_NAMES.CLEANUP_AUDIT_LOGS,
);
expect(mockRepository.createAuditLog).toHaveBeenCalledWith(
expect.objectContaining({
actorId: "admin-1",
action: AUDIT_ACTIONS.CRON_JOB_TOGGLED,
details: {
jobName: CRON_JOB_NAMES.CLEANUP_AUDIT_LOGS,
isEnabled: false,
},
}),
);
expect(result.isEnabled).toBe(false);
});
it("should enable scheduler when enabled is true", async () => {
await cronService.toggleJob(CRON_JOB_NAMES.CLEANUP_AUDIT_LOGS, true);
expect(mockQueueService.enableJobScheduler).toHaveBeenCalledWith(
CRON_JOB_NAMES.CLEANUP_AUDIT_LOGS,
);
});
it("should throw 404 when job does not exist", async () => {
await expect(
cronService.toggleJob("non-existent-job" as any, true),
).rejects.toThrow();
});
});
describe("triggerJob / executeJob", () => {
it("should execute cleanup-audit-logs and record audit log", async () => {
const result = await cronService.triggerJob(
CRON_JOB_NAMES.CLEANUP_AUDIT_LOGS,
{ retentionDays: 15 },
{ actorId: "admin-1" },
);
expect(result.success).toBe(true);
expect(mockRepository.cleanupAuditLogs).toHaveBeenCalled();
expect(result.data).toEqual(
expect.objectContaining({
deletedCount: 45,
retentionDays: 15,
}),
);
expect(mockRepository.createAuditLog).toHaveBeenCalledWith(
expect.objectContaining({
action: AUDIT_ACTIONS.CRON_JOB_TRIGGERED,
}),
);
});
it("should execute cleanup-unconfirmed-uploads and delete expired export files", async () => {
const result = await cronService.triggerJob(
CRON_JOB_NAMES.CLEANUP_UNCONFIRMED_UPLOADS,
);
expect(result.success).toBe(true);
expect(mockRepository.cleanupExpiredExports).toHaveBeenCalled();
expect(mockStorageService.deleteFile).toHaveBeenCalledWith("exports/test1.zip");
expect(mockStorageService.deleteFile).toHaveBeenCalledWith("exports/test2.zip");
expect(result.data).toEqual(
expect.objectContaining({
cleanedExportsCount: 2,
}),
);
});
it("should execute cleanup-expired-tokens", async () => {
const result = await cronService.triggerJob(
CRON_JOB_NAMES.CLEANUP_EXPIRED_TOKENS,
);
expect(result.success).toBe(true);
expect(mockRepository.cleanupExpiredTokens).toHaveBeenCalled();
expect(result.data).toEqual({ deletedRefreshTokens: 12 });
});
it("should execute daily-summary-digest and send email to admins", async () => {
const result = await cronService.triggerJob(
CRON_JOB_NAMES.DAILY_SUMMARY_DIGEST,
);
expect(result.success).toBe(true);
expect(mockRepository.getDigestStats).toHaveBeenCalled();
expect(mockRepository.getAdminEmails).toHaveBeenCalled();
expect(mockMailService.sendDigestEmail).toHaveBeenCalledWith(
["admin@datacrawler.com"],
expect.stringContaining("Báo Cáo Hoạt Động Hàng Ngày"),
expect.objectContaining({ period: "DAILY" }),
);
});
it("should execute weekly-summary-digest and send email to admins", async () => {
const result = await cronService.triggerJob(
CRON_JOB_NAMES.WEEKLY_SUMMARY_DIGEST,
);
expect(result.success).toBe(true);
expect(mockRepository.getDigestStats).toHaveBeenCalled();
expect(mockMailService.sendDigestEmail).toHaveBeenCalledWith(
["admin@datacrawler.com"],
expect.stringContaining("Báo Cáo Hiệu Suất Hệ Thống Hàng Tuần"),
expect.objectContaining({ period: "WEEKLY" }),
);
});
});
});
import { Request, Response, NextFunction } from "express";
import { CronService, cronService } from "./cron.service";
import { CronJobName } from "../../common/constants/cron.constant";
export class CronController {
constructor(private readonly service: CronService = cronService) {}
/**
* GET /api/v1/cron/jobs - Lấy danh sách toàn bộ tác vụ định kỳ
*/
listJobs = async (
req: Request,
res: Response,
next: NextFunction,
): Promise<void> => {
try {
const { search } = req.query as { search?: string };
const data = await this.service.listJobs(search);
res.status(200).json({
success: true,
data,
});
} catch (error) {
next(error);
}
};
/**
* POST /api/v1/cron/jobs/:jobName/trigger - Kích hoạt chạy ngay một tác vụ
*/
triggerJob = async (
req: Request,
res: Response,
next: NextFunction,
): Promise<void> => {
try {
const { jobName } = req.params as { jobName: CronJobName };
const { params } = req.body as { params?: Record<string, unknown> };
const actorContext = {
actorId: req.user?.id,
ipAddress: req.ip || req.socket.remoteAddress,
userAgent: req.get("user-agent"),
};
const data = await this.service.triggerJob(jobName, params, actorContext);
res.status(200).json({
success: true,
data,
});
} catch (error) {
next(error);
}
};
/**
* PATCH /api/v1/cron/jobs/:jobName/toggle - Bật hoặc tắt lịch chạy tự động của tác vụ
*/
toggleJob = async (
req: Request,
res: Response,
next: NextFunction,
): Promise<void> => {
try {
const { jobName } = req.params as { jobName: CronJobName };
const { enabled } = req.body as { enabled: boolean };
const actorContext = {
actorId: req.user?.id,
ipAddress: req.ip || req.socket.remoteAddress,
userAgent: req.get("user-agent"),
};
const data = await this.service.toggleJob(jobName, enabled, actorContext);
res.status(200).json({
success: true,
data,
});
} catch (error) {
next(error);
}
};
}
export const cronController = new CronController();
import { CronJobName, CronJobStatus } from "../../common/constants/cron.constant";
export interface CronJobItemDto {
name: CronJobName;
cron: string;
description: string;
isEnabled: boolean;
lastRun?: string;
lastStatus?: CronJobStatus;
lastDurationMs?: number;
}
export interface CronJobExecutionResultDto {
jobName: CronJobName;
success: boolean;
durationMs: number;
data?: Record<string, unknown>;
error?: string;
}
export interface TriggerJobBodyDto {
params?: Record<string, unknown>;
}
export interface ToggleJobBodyDto {
enabled: boolean;
}
export interface ListJobsQueryDto {
search?: string;
}
export interface AuditContext {
actorId?: string;
ipAddress?: string;
userAgent?: string;
source?: "MANUAL_TRIGGER" | "SCHEDULER";
}
import { prisma } from "../../database/prisma.client";
import { Prisma } from "@prisma/client";
import {
CRON_JOB_STATUS,
CRON_SYSTEM_CONFIG_KEY,
CronJobStatus,
} from "../../common/constants/cron.constant";
import { AUDIT_ACTIONS } from "../../common/constants/audit-action.constant";
export interface CreateCronAuditLogInput {
actorId?: string;
action: string;
details?: Record<string, unknown>;
ipAddress?: string;
userAgent?: string;
}
export interface JobRunHistoryItem {
lastRun: Date;
status: CronJobStatus;
durationMs?: number;
}
export class CronRepository {
/**
* Lấy cấu hình cờ trạng thái bật/tắt của toàn bộ các Cron Jobs từ SystemConfig
*/
async getJobStatuses(): Promise<Record<string, boolean>> {
const config = await prisma.systemConfig.findUnique({
where: { key: CRON_SYSTEM_CONFIG_KEY },
});
if (!config || typeof config.value !== "object" || config.value === null) {
return {};
}
return config.value as Record<string, boolean>;
}
/**
* Cập nhật trạng thái bật/tắt của một Cron Job cụ thể trong SystemConfig
*/
async setJobStatus(
jobName: string,
enabled: boolean,
): Promise<Record<string, boolean>> {
const currentStatuses = await this.getJobStatuses();
const updatedStatuses = {
...currentStatuses,
[jobName]: enabled,
};
await prisma.systemConfig.upsert({
where: { key: CRON_SYSTEM_CONFIG_KEY },
update: {
value: updatedStatuses as unknown as Prisma.InputJsonValue,
},
create: {
key: CRON_SYSTEM_CONFIG_KEY,
value: updatedStatuses as unknown as Prisma.InputJsonValue,
description: "Bảng trạng thái kích hoạt tự động của các tác vụ nền định kỳ (Cron Jobs)",
category: "FEATURE_FLAG",
isPublic: false,
},
});
return updatedStatuses;
}
/**
* Lấy lịch sử thực thi gần nhất của danh sách các tác vụ từ AuditLog
*/
async getLatestRunsForJobs(
jobNames: string[],
): Promise<Map<string, JobRunHistoryItem>> {
const result = new Map<string, JobRunHistoryItem>();
if (jobNames.length === 0) return result;
const recentLogs = await prisma.auditLog.findMany({
where: {
action: {
in: [
AUDIT_ACTIONS.CRON_JOB_EXECUTED,
AUDIT_ACTIONS.CRON_JOB_TRIGGERED,
],
},
},
orderBy: { createdAt: "desc" },
take: 100,
});
for (const log of recentLogs) {
const details = log.details as Record<string, unknown> | null;
const jobName = details?.jobName as string | undefined;
if (jobName && jobNames.includes(jobName) && !result.has(jobName)) {
const isSuccess = details?.success !== false;
result.set(jobName, {
lastRun: log.createdAt,
status: isSuccess ? CRON_JOB_STATUS.SUCCESS : CRON_JOB_STATUS.FAILED,
durationMs: typeof details?.durationMs === "number" ? details.durationMs : undefined,
});
}
}
return result;
}
/**
* Ghi vết kiểm toán cho tác vụ định kỳ
*/
async createAuditLog(data: CreateCronAuditLogInput) {
return prisma.auditLog.create({
data: {
userId: data.actorId || null,
action: data.action,
details: (data.details ?? undefined) as Prisma.InputJsonValue | undefined,
ipAddress: data.ipAddress || null,
userAgent: data.userAgent || null,
},
});
}
/**
* Xóa vĩnh viễn các bản ghi AuditLog cũ hơn thời điểm chỉ định
*/
async cleanupAuditLogs(cutoffDate: Date): Promise<number> {
const result = await prisma.auditLog.deleteMany({
where: {
createdAt: {
lt: cutoffDate,
},
},
});
return result.count;
}
/**
* Dọn dẹp RefreshToken đã hết hạn trong 1 database transaction
*/
async cleanupExpiredTokens(now: Date): Promise<{ deletedRefreshTokens: number }> {
const [refreshTokensResult] = await prisma.$transaction([
prisma.refreshToken.deleteMany({
where: {
expiresAt: {
lt: now,
},
},
}),
]);
return {
deletedRefreshTokens: refreshTokensResult.count,
};
}
/**
* Dọn dẹp các tệp xuất dữ liệu CrawlExport đã hết hạn (expiredAt < now)
*/
async cleanupExpiredExports(
now: Date,
): Promise<{ deletedCount: number; filePaths: string[] }> {
const expiredExports = await prisma.crawlExport.findMany({
where: {
expiredAt: {
not: null,
lt: now,
},
},
select: {
id: true,
filePath: true,
},
});
if (expiredExports.length === 0) {
return { deletedCount: 0, filePaths: [] };
}
const ids = expiredExports.map((e) => e.id);
const filePaths = expiredExports.map((e) => e.filePath).filter(Boolean);
await prisma.crawlExport.deleteMany({
where: {
id: {
in: ids,
},
},
});
return {
deletedCount: ids.length,
filePaths,
};
}
/**
* Lấy danh sách đường dẫn avatar đang được người dùng sử dụng
*/
async getActiveAvatarUrls(): Promise<string[]> {
const users = await prisma.user.findMany({
where: {
avatarUrl: {
not: null,
},
},
select: {
avatarUrl: true,
},
});
return users
.map((u) => u.avatarUrl)
.filter((url): url is string => Boolean(url));
}
/**
* Tổng hợp các chỉ số hoạt động trong một khoảng thời gian
*/
async getDigestStats(startDate: Date, endDate: Date) {
const [
newUsers,
crawlJobsTotal,
crawlJobsCompleted,
crawlJobsFailed,
crawledPages,
exportsGenerated,
webhookDeliveries,
auditLogsRecorded,
] = await Promise.all([
prisma.user.count({
where: { createdAt: { gte: startDate, lte: endDate } },
}),
prisma.crawlJob.count({
where: { createdAt: { gte: startDate, lte: endDate } },
}),
prisma.crawlJob.count({
where: {
createdAt: { gte: startDate, lte: endDate },
status: "COMPLETED",
},
}),
prisma.crawlJob.count({
where: {
createdAt: { gte: startDate, lte: endDate },
status: "FAILED",
},
}),
prisma.crawlPage.count({
where: {
crawledAt: { gte: startDate, lte: endDate },
},
}),
prisma.crawlExport.count({
where: { createdAt: { gte: startDate, lte: endDate } },
}),
prisma.webhookDelivery.count({
where: { createdAt: { gte: startDate, lte: endDate } },
}),
prisma.auditLog.count({
where: { createdAt: { gte: startDate, lte: endDate } },
}),
]);
return {
newUsers,
crawlJobsTotal,
crawlJobsCompleted,
crawlJobsFailed,
crawledPages,
exportsGenerated,
webhookDeliveries,
auditLogsRecorded,
};
}
/**
* Lấy danh sách email quản trị viên hệ thống để gửi báo cáo
*/
async getAdminEmails(): Promise<string[]> {
const admins = await prisma.user.findMany({
where: {
role: "ADMIN",
isActive: true,
},
select: {
email: true,
},
});
return admins.map((a) => a.email);
}
}
export const cronRepository = new CronRepository();
import { Router } from "express";
import { cronController } from "./cron.controller";
import { authMiddleware } from "../../middlewares/auth.middleware";
import { requirePermission } from "../../middlewares/permission.middleware";
import {
validate,
validateParams,
validateQuery,
} from "../../middlewares/validate.middleware";
import {
cronJobParamsSchema,
listCronJobsQuerySchema,
toggleCronJobBodySchema,
triggerCronJobBodySchema,
} from "./cron.validation";
import { PERMISSIONS } from "../../common/constants/permission.constant";
const router = Router();
/**
* 1. GET /api/v1/cron/jobs
* Danh sách toàn bộ các tác vụ định kỳ và lịch biểu
*/
router.get(
"/jobs",
authMiddleware,
requirePermission(PERMISSIONS.CRON_JOB_READ),
validateQuery(listCronJobsQuerySchema),
cronController.listJobs,
);
/**
* 2. POST /api/v1/cron/jobs/:jobName/trigger
* Kích hoạt chạy ngay một tác vụ thủ công
*/
router.post(
"/jobs/:jobName/trigger",
authMiddleware,
requirePermission(PERMISSIONS.CRON_JOB_MANAGE),
validateParams(cronJobParamsSchema),
validate(triggerCronJobBodySchema),
cronController.triggerJob,
);
/**
* 3. PATCH /api/v1/cron/jobs/:jobName/toggle
* Bật hoặc tắt kích hoạt tự động theo lịch của một tác vụ
*/
router.patch(
"/jobs/:jobName/toggle",
authMiddleware,
requirePermission(PERMISSIONS.CRON_JOB_MANAGE),
validateParams(cronJobParamsSchema),
validate(toggleCronJobBodySchema),
cronController.toggleJob,
);
export default router;
This diff is collapsed.
import { z } from "zod";
import { CRON_JOB_NAMES } from "../../common/constants/cron.constant";
export const cronJobParamsSchema = z.object({
jobName: z.nativeEnum(CRON_JOB_NAMES, {
errorMap: () => ({ message: "Tên tác vụ không tồn tại trong hệ thống" }),
}),
});
export const triggerCronJobBodySchema = z.object({
params: z.record(z.unknown()).optional(),
});
export const toggleCronJobBodySchema = z.object({
enabled: z.boolean({
required_error: "Trường 'enabled' là bắt buộc (true hoặc false)",
}),
});
export const listCronJobsQuerySchema = z.object({
search: z.string().trim().optional(),
});
...@@ -120,4 +120,100 @@ export class MailService { ...@@ -120,4 +120,100 @@ export class MailService {
await this.transporter.sendMail(mailOptions); await this.transporter.sendMail(mailOptions);
} }
async sendDigestEmail(
recipients: string[],
subject: string,
digest: {
period: "DAILY" | "WEEKLY";
startDate: Date;
endDate: Date;
stats: {
newUsers: number;
crawlJobsTotal: number;
crawlJobsCompleted: number;
crawlJobsFailed: number;
crawledPages: number;
exportsGenerated: number;
webhookDeliveries: number;
auditLogsRecorded: number;
};
},
): Promise<void> {
if (!recipients.length) return;
const periodLabel = digest.period === "DAILY" ? "Hàng Ngày" : "Hàng Tuần";
const dateRange = `${digest.startDate.toLocaleDateString("vi-VN")} - ${digest.endDate.toLocaleDateString("vi-VN")}`;
const mailOptions = {
from: mailConfig.from,
to: recipients.join(", "),
subject,
html: `
<div style="font-family: -apple-system, BlinkMacSystemFont, 'Segoe UI', Roboto, Helvetica, Arial, sans-serif; max-width: 620px; margin: 0 auto; padding: 24px; background-color: #fcfdfc; border: 1px solid #d1fae5; border-radius: 16px;">
<div style="text-align: center; margin-bottom: 24px;">
<div style="display: inline-block; padding: 6px 14px; background-color: #ecfdf5; border: 1px solid #a7f3d0; border-radius: 9999px; font-size: 12px; font-weight: 600; color: #047857;">
🌿 Báo Cáo Tổng Hợp ${periodLabel}
</div>
<h2 style="color: #064e3b; margin: 12px 0 4px 0; font-size: 22px;">Hệ Thống Data Crawler</h2>
<p style="color: #6b7280; font-size: 13px; margin: 0;">Khung thời gian: ${dateRange}</p>
</div>
<div style="background-color: #ffffff; border: 1px solid #e5e7eb; border-radius: 12px; padding: 18px; margin-bottom: 20px;">
<h3 style="color: #111827; font-size: 15px; margin-top: 0; margin-bottom: 12px; border-bottom: 1px solid #f3f4f6; padding-bottom: 8px;">
📊 Chỉ Số Hoạt Động Cốt Lõi
</h3>
<table style="width: 100%; border-collapse: collapse; font-size: 14px;">
<tr>
<td style="padding: 8px 0; color: #4b5563;">Người dùng mới đăng ký:</td>
<td style="padding: 8px 0; text-align: right; font-weight: bold; color: #065f46;">${digest.stats.newUsers}</td>
</tr>
<tr>
<td style="padding: 8px 0; color: #4b5563;">Tổng tác vụ cào dữ liệu:</td>
<td style="padding: 8px 0; text-align: right; font-weight: bold; color: #111827;">${digest.stats.crawlJobsTotal}</td>
</tr>
<tr>
<td style="padding: 8px 0; color: #4b5563;">Tác vụ hoàn thành thành công:</td>
<td style="padding: 8px 0; text-align: right; font-weight: bold; color: #059669;">${digest.stats.crawlJobsCompleted}</td>
</tr>
<tr>
<td style="padding: 8px 0; color: #4b5563;">Tác vụ thất bại:</td>
<td style="padding: 8px 0; text-align: right; font-weight: bold; color: ${digest.stats.crawlJobsFailed > 0 ? "#dc2626" : "#4b5563"};">${digest.stats.crawlJobsFailed}</td>
</tr>
<tr>
<td style="padding: 8px 0; color: #4b5563;">Tổng số trang web đã crawl:</td>
<td style="padding: 8px 0; text-align: right; font-weight: bold; color: #111827;">${digest.stats.crawledPages}</td>
</tr>
<tr>
<td style="padding: 8px 0; color: #4b5563;">Tệp dữ liệu xuất bản (Exports):</td>
<td style="padding: 8px 0; text-align: right; font-weight: bold; color: #111827;">${digest.stats.exportsGenerated}</td>
</tr>
<tr>
<td style="padding: 8px 0; color: #4b5563;">Thông báo Webhook đã gửi:</td>
<td style="padding: 8px 0; text-align: right; font-weight: bold; color: #111827;">${digest.stats.webhookDeliveries}</td>
</tr>
<tr>
<td style="padding: 8px 0; color: #4b5563;">Nhật ký kiểm toán ghi nhận:</td>
<td style="padding: 8px 0; text-align: right; font-weight: bold; color: #111827;">${digest.stats.auditLogsRecorded}</td>
</tr>
</table>
</div>
<div style="text-align: center; margin: 24px 0 12px 0;">
<a href="${mailConfig.frontendUrl}/dashboard" style="background-color: #059669; color: #ffffff; padding: 10px 20px; text-decoration: none; border-radius: 8px; font-weight: 600; font-size: 13px; display: inline-block;">
Mở Bảng Điều Khiển Quản Trị
</a>
</div>
<hr style="border: 0; border-top: 1px solid #e5e7eb; margin: 20px 0;">
<p style="color: #9ca3af; font-size: 12px; text-align: center; margin: 0;">
Đây là email tự động từ Phân hệ Tác vụ Định kỳ Data Crawler. Bạn nhận được thư này vì có quyền Quản trị viên hệ thống.
</p>
</div>
`,
};
await this.transporter.sendMail(mailOptions);
}
} }
...@@ -110,4 +110,16 @@ describe("PermissionService", () => { ...@@ -110,4 +110,16 @@ describe("PermissionService", () => {
expect(repository.findUserRoleSlugs).toHaveBeenCalledTimes(1); expect(repository.findUserRoleSlugs).toHaveBeenCalledTimes(1);
}); });
}); });
describe("ensureSystemPermissions", () => {
it("should call repository.ensureSystemPermissions and invalidate cache", async () => {
repository.ensureSystemPermissions = jest.fn().mockResolvedValue(undefined);
const invalidateSpy = jest.spyOn(authorizationCache, "invalidateAll");
await service.ensureSystemPermissions();
expect(repository.ensureSystemPermissions).toHaveBeenCalledTimes(1);
expect(invalidateSpy).toHaveBeenCalledTimes(1);
});
});
}); });
import { prisma } from "../../database/prisma.client"; import { prisma } from "../../database/prisma.client";
import { Permission } from "../../common/types/database.types"; import { Permission } from "../../common/types/database.types";
import {
SYSTEM_PERMISSIONS_CATALOG,
SYSTEM_ROLE_DEFAULT_PERMISSIONS,
} from "../../common/constants/permission.constant";
import {
SYSTEM_ROLE_SLUGS,
SYSTEM_ROLES_METADATA,
SystemRoleSlug,
} from "../../common/constants/system-role.constant";
export class PermissionRepository { export class PermissionRepository {
async findAll(query?: { async findAll(query?: {
...@@ -90,4 +99,91 @@ export class PermissionRepository { ...@@ -90,4 +99,91 @@ export class PermissionRepository {
return assignments.map((a) => a.role.slug); return assignments.map((a) => a.role.slug);
} }
async ensureSystemPermissions(): Promise<void> {
// 1. Upsert all permissions from SYSTEM_PERMISSIONS_CATALOG
for (const perm of SYSTEM_PERMISSIONS_CATALOG) {
await prisma.permission.upsert({
where: { slug: perm.slug },
update: {
name: perm.name,
description: perm.description,
resource: perm.resource,
action: perm.action,
isSystem: perm.isSystem,
},
create: {
name: perm.name,
slug: perm.slug,
description: perm.description,
resource: perm.resource,
action: perm.action,
isSystem: perm.isSystem,
},
});
}
// 2. Ensure system roles exist
for (const slug of Object.values(SYSTEM_ROLE_SLUGS)) {
const meta = SYSTEM_ROLES_METADATA[slug as SystemRoleSlug];
if (meta) {
await prisma.role.upsert({
where: { slug },
update: {
name: meta.name,
description: meta.description,
isSystem: meta.isSystem,
isActive: true,
},
create: {
name: meta.name,
slug: meta.slug,
description: meta.description,
isSystem: meta.isSystem,
isActive: true,
},
});
}
}
// 3. Fetch all permissions map: slug -> id
const allPermissions = await prisma.permission.findMany({
select: { id: true, slug: true },
});
const permissionMap = new Map<string, string>(
allPermissions.map((p) => [p.slug, p.id]),
);
// 4. For each system role, ensure default permissions are assigned in role_permissions
for (const slug of Object.values(SYSTEM_ROLE_SLUGS)) {
const role = await prisma.role.findUnique({
where: { slug },
select: { id: true },
});
if (!role) continue;
const defaultPermSlugs =
SYSTEM_ROLE_DEFAULT_PERMISSIONS[slug as SystemRoleSlug] || [];
for (const permSlug of defaultPermSlugs) {
const permId = permissionMap.get(permSlug);
if (permId) {
await prisma.rolePermission.upsert({
where: {
roleId_permissionId: {
roleId: role.id,
permissionId: permId,
},
},
update: {},
create: {
roleId: role.id,
permissionId: permId,
},
});
}
}
}
}
} }
...@@ -66,4 +66,11 @@ export class PermissionService { ...@@ -66,4 +66,11 @@ export class PermissionService {
authorizationCache.setCachedRoles(userId, roles); authorizationCache.setCachedRoles(userId, roles);
return roles; return roles;
} }
async ensureSystemPermissions(): Promise<void> {
await this.repository.ensureSystemPermissions();
authorizationCache.invalidateAll();
}
} }
export const permissionService = new PermissionService();
...@@ -279,4 +279,42 @@ describe("RoleService", () => { ...@@ -279,4 +279,42 @@ describe("RoleService", () => {
); );
}); });
}); });
describe("resetRoleQuota", () => {
it("should reset role quota for all assigned users and log audit", async () => {
repository.findById.mockResolvedValue({
id: "r-cust",
slug: "custom_role",
} as any);
repository.resetRoleQuota.mockResolvedValue({ count: 5 } as any);
const result = await service.resetRoleQuota("r-cust", true, {
actorId: "actor-admin",
});
expect(result).toEqual({ affectedUsers: 5 });
expect(repository.resetRoleQuota).toHaveBeenCalledWith("r-cust", true);
expect(auditLogService.log).toHaveBeenCalledWith(
expect.objectContaining({
userId: "actor-admin",
action: AUDIT_ACTIONS.ROLE_QUOTA_RESET,
details: expect.objectContaining({
roleId: "r-cust",
roleSlug: "custom_role",
syncLimits: true,
affectedUsers: 5,
}),
}),
);
});
it("should throw 404 if role does not exist", async () => {
repository.findById.mockResolvedValue(null);
await expect(service.resetRoleQuota("non-existent")).rejects.toThrow(
AppError,
);
});
});
}); });
...@@ -110,6 +110,33 @@ export class RoleController { ...@@ -110,6 +110,33 @@ export class RoleController {
} }
}; };
resetRoleQuota = async (
req: Request,
res: Response,
next: NextFunction,
): Promise<void> => {
try {
const syncLimits = !!req.body?.syncLimits;
const result = await this.service.resetRoleQuota(
req.params.id,
syncLimits,
{
actorId: req.user.id,
ipAddress: req.ip,
userAgent: req.headers["user-agent"] as string,
},
);
res.json({
success: true,
message: "Role quota reset successfully",
data: result,
});
} catch (error) {
next(error);
}
};
getRolePermissions = async ( getRolePermissions = async (
req: Request, req: Request,
res: Response, res: Response,
......
...@@ -3,12 +3,27 @@ export interface CreateRoleDto { ...@@ -3,12 +3,27 @@ export interface CreateRoleDto {
slug: string; slug: string;
description?: string; description?: string;
permissionIds?: string[]; permissionIds?: string[];
maxPagesLimit?: number;
maxJobsPerDayLimit?: number;
maxConcurrentJobsLimit?: number;
maxPagesPerMonthLimit?: number | null;
maxJobsPerMonthLimit?: number | null;
} }
export interface UpdateRoleDto { export interface UpdateRoleDto {
name?: string; name?: string;
description?: string; description?: string;
isActive?: boolean; isActive?: boolean;
maxPagesLimit?: number;
maxJobsPerDayLimit?: number;
maxConcurrentJobsLimit?: number;
maxPagesPerMonthLimit?: number | null;
maxJobsPerMonthLimit?: number | null;
syncUsersQuota?: boolean;
}
export interface ResetRoleQuotaDto {
syncLimits?: boolean;
} }
export interface RoleQueryDto { export interface RoleQueryDto {
...@@ -38,6 +53,11 @@ export interface RoleResponseDto { ...@@ -38,6 +53,11 @@ export interface RoleResponseDto {
description: string | null; description: string | null;
isSystem: boolean; isSystem: boolean;
isActive: boolean; isActive: boolean;
maxPagesLimit: number;
maxJobsPerDayLimit: number;
maxConcurrentJobsLimit: number;
maxPagesPerMonthLimit: number | null;
maxJobsPerMonthLimit: number | null;
createdAt: Date; createdAt: Date;
updatedAt: Date; updatedAt: Date;
permissions?: RolePermissionItemDto[]; permissions?: RolePermissionItemDto[];
......
...@@ -114,6 +114,11 @@ export class RoleRepository { ...@@ -114,6 +114,11 @@ export class RoleRepository {
description?: string; description?: string;
isSystem?: boolean; isSystem?: boolean;
permissionIds?: string[]; permissionIds?: string[];
maxPagesLimit?: number;
maxJobsPerDayLimit?: number;
maxConcurrentJobsLimit?: number;
maxPagesPerMonthLimit?: number | null;
maxJobsPerMonthLimit?: number | null;
}) { }) {
return prisma.$transaction(async (tx) => { return prisma.$transaction(async (tx) => {
const role = await tx.role.create({ const role = await tx.role.create({
...@@ -123,6 +128,11 @@ export class RoleRepository { ...@@ -123,6 +128,11 @@ export class RoleRepository {
description: data.description, description: data.description,
isSystem: data.isSystem ?? false, isSystem: data.isSystem ?? false,
isActive: true, isActive: true,
...(data.maxPagesLimit !== undefined && { maxPagesLimit: data.maxPagesLimit }),
...(data.maxJobsPerDayLimit !== undefined && { maxJobsPerDayLimit: data.maxJobsPerDayLimit }),
...(data.maxConcurrentJobsLimit !== undefined && { maxConcurrentJobsLimit: data.maxConcurrentJobsLimit }),
...(data.maxPagesPerMonthLimit !== undefined && { maxPagesPerMonthLimit: data.maxPagesPerMonthLimit }),
...(data.maxJobsPerMonthLimit !== undefined && { maxJobsPerMonthLimit: data.maxJobsPerMonthLimit }),
}, },
}); });
...@@ -156,6 +166,11 @@ export class RoleRepository { ...@@ -156,6 +166,11 @@ export class RoleRepository {
name?: string; name?: string;
description?: string; description?: string;
isActive?: boolean; isActive?: boolean;
maxPagesLimit?: number;
maxJobsPerDayLimit?: number;
maxConcurrentJobsLimit?: number;
maxPagesPerMonthLimit?: number | null;
maxJobsPerMonthLimit?: number | null;
}, },
) { ) {
return prisma.role.update({ return prisma.role.update({
...@@ -172,6 +187,72 @@ export class RoleRepository { ...@@ -172,6 +187,72 @@ export class RoleRepository {
}); });
} }
async syncUsersQuota(
roleId: string,
limits: {
maxPagesLimit?: number;
maxJobsPerDayLimit?: number;
maxConcurrentJobsLimit?: number;
maxPagesPerMonthLimit?: number | null;
maxJobsPerMonthLimit?: number | null;
},
) {
const updateData: Prisma.UserUpdateManyMutationInput = {};
if (limits.maxPagesLimit !== undefined) updateData.maxPagesLimit = limits.maxPagesLimit;
if (limits.maxJobsPerDayLimit !== undefined) updateData.maxJobsPerDayLimit = limits.maxJobsPerDayLimit;
if (limits.maxConcurrentJobsLimit !== undefined) updateData.maxConcurrentJobsLimit = limits.maxConcurrentJobsLimit;
if (limits.maxPagesPerMonthLimit !== undefined) updateData.maxPagesPerMonthLimit = limits.maxPagesPerMonthLimit;
if (limits.maxJobsPerMonthLimit !== undefined) updateData.maxJobsPerMonthLimit = limits.maxJobsPerMonthLimit;
if (Object.keys(updateData).length === 0) return { count: 0 };
return prisma.user.updateMany({
where: {
userRoles: {
some: { roleId },
},
},
data: updateData,
});
}
async resetRoleQuota(roleId: string, syncLimits: boolean = false) {
const now = new Date();
const updateData: Prisma.UserUpdateManyMutationInput = {
quotaResetAt: now,
};
if (syncLimits) {
const role = await prisma.role.findUnique({
where: { id: roleId },
select: {
maxPagesLimit: true,
maxJobsPerDayLimit: true,
maxConcurrentJobsLimit: true,
maxPagesPerMonthLimit: true,
maxJobsPerMonthLimit: true,
},
});
if (role) {
updateData.maxPagesLimit = role.maxPagesLimit;
updateData.maxJobsPerDayLimit = role.maxJobsPerDayLimit;
updateData.maxConcurrentJobsLimit = role.maxConcurrentJobsLimit;
updateData.maxPagesPerMonthLimit = role.maxPagesPerMonthLimit;
updateData.maxJobsPerMonthLimit = role.maxJobsPerMonthLimit;
}
}
return prisma.user.updateMany({
where: {
userRoles: {
some: { roleId },
},
},
data: updateData,
});
}
async delete(id: string) { async delete(id: string) {
return prisma.role.delete({ return prisma.role.delete({
where: { id }, where: { id },
......
...@@ -10,6 +10,7 @@ import { ...@@ -10,6 +10,7 @@ import {
import { import {
createRoleSchema, createRoleSchema,
updateRoleSchema, updateRoleSchema,
resetRoleQuotaSchema,
assignRolePermissionsSchema, assignRolePermissionsSchema,
listRolesQuerySchema, listRolesQuerySchema,
roleParamsSchema, roleParamsSchema,
...@@ -60,6 +61,15 @@ router.delete( ...@@ -60,6 +61,15 @@ router.delete(
controller.delete, controller.delete,
); );
router.post(
"/:id/reset-quota",
authMiddleware,
requirePermission(PERMISSIONS.ROLES_UPDATE),
validateParams(roleParamsSchema),
validate(resetRoleQuotaSchema),
controller.resetRoleQuota,
);
router.get( router.get(
"/:id/permissions", "/:id/permissions",
authMiddleware, authMiddleware,
......
...@@ -25,6 +25,11 @@ interface RoleWithPermissions { ...@@ -25,6 +25,11 @@ interface RoleWithPermissions {
description: string | null; description: string | null;
isSystem: boolean; isSystem: boolean;
isActive: boolean; isActive: boolean;
maxPagesLimit: number;
maxJobsPerDayLimit: number;
maxConcurrentJobsLimit: number;
maxPagesPerMonthLimit: number | null;
maxJobsPerMonthLimit: number | null;
createdAt: Date; createdAt: Date;
updatedAt: Date; updatedAt: Date;
rolePermissions?: Array<{ rolePermissions?: Array<{
...@@ -53,6 +58,11 @@ export class RoleService { ...@@ -53,6 +58,11 @@ export class RoleService {
description: role.description ?? null, description: role.description ?? null,
isSystem: role.isSystem, isSystem: role.isSystem,
isActive: role.isActive, isActive: role.isActive,
maxPagesLimit: role.maxPagesLimit,
maxJobsPerDayLimit: role.maxJobsPerDayLimit,
maxConcurrentJobsLimit: role.maxConcurrentJobsLimit,
maxPagesPerMonthLimit: role.maxPagesPerMonthLimit,
maxJobsPerMonthLimit: role.maxJobsPerMonthLimit,
createdAt: role.createdAt, createdAt: role.createdAt,
updatedAt: role.updatedAt, updatedAt: role.updatedAt,
permissions: role.rolePermissions permissions: role.rolePermissions
...@@ -121,6 +131,11 @@ export class RoleService { ...@@ -121,6 +131,11 @@ export class RoleService {
description: dto.description, description: dto.description,
isSystem: false, isSystem: false,
permissionIds: dto.permissionIds, permissionIds: dto.permissionIds,
maxPagesLimit: dto.maxPagesLimit,
maxJobsPerDayLimit: dto.maxJobsPerDayLimit,
maxConcurrentJobsLimit: dto.maxConcurrentJobsLimit,
maxPagesPerMonthLimit: dto.maxPagesPerMonthLimit,
maxJobsPerMonthLimit: dto.maxJobsPerMonthLimit,
}); });
if (context?.actorId) { if (context?.actorId) {
...@@ -165,8 +180,23 @@ export class RoleService { ...@@ -165,8 +180,23 @@ export class RoleService {
name: dto.name, name: dto.name,
description: dto.description, description: dto.description,
isActive: dto.isActive, isActive: dto.isActive,
maxPagesLimit: dto.maxPagesLimit,
maxJobsPerDayLimit: dto.maxJobsPerDayLimit,
maxConcurrentJobsLimit: dto.maxConcurrentJobsLimit,
maxPagesPerMonthLimit: dto.maxPagesPerMonthLimit,
maxJobsPerMonthLimit: dto.maxJobsPerMonthLimit,
}); });
if (dto.syncUsersQuota) {
await this.repository.syncUsersQuota(id, {
maxPagesLimit: updated.maxPagesLimit,
maxJobsPerDayLimit: updated.maxJobsPerDayLimit,
maxConcurrentJobsLimit: updated.maxConcurrentJobsLimit,
maxPagesPerMonthLimit: updated.maxPagesPerMonthLimit,
maxJobsPerMonthLimit: updated.maxJobsPerMonthLimit,
});
}
if (context?.actorId) { if (context?.actorId) {
await this.auditLogService.log({ await this.auditLogService.log({
userId: context.actorId, userId: context.actorId,
...@@ -185,6 +215,36 @@ export class RoleService { ...@@ -185,6 +215,36 @@ export class RoleService {
return this.formatRole(updated); return this.formatRole(updated);
} }
async resetRoleQuota(
id: string,
syncLimits: boolean = false,
context?: AuditContext,
): Promise<{ affectedUsers: number }> {
const existing = await this.repository.findById(id);
if (!existing) {
throw new AppError("Role not found", 404, ERROR_CODE.ROLE_NOT_FOUND);
}
const result = await this.repository.resetRoleQuota(id, syncLimits);
if (context?.actorId) {
await this.auditLogService.log({
userId: context.actorId,
action: AUDIT_ACTIONS.ROLE_QUOTA_RESET,
ipAddress: context.ipAddress,
userAgent: context.userAgent,
details: {
roleId: id,
roleSlug: existing.slug,
syncLimits,
affectedUsers: result.count,
},
});
}
return { affectedUsers: result.count };
}
async delete(id: string, context?: AuditContext): Promise<void> { async delete(id: string, context?: AuditContext): Promise<void> {
const existing = await this.repository.findById(id); const existing = await this.repository.findById(id);
if (!existing) { if (!existing) {
......
...@@ -19,6 +19,11 @@ export const createRoleSchema = z.object({ ...@@ -19,6 +19,11 @@ export const createRoleSchema = z.object({
permissionIds: z permissionIds: z
.array(z.string().uuid("Invalid permission ID format")) .array(z.string().uuid("Invalid permission ID format"))
.optional(), .optional(),
maxPagesLimit: z.number().int().min(1).max(100000).optional(),
maxJobsPerDayLimit: z.number().int().min(1).max(10000).optional(),
maxConcurrentJobsLimit: z.number().int().min(1).max(100).optional(),
maxPagesPerMonthLimit: z.number().int().min(1).max(1000000).nullable().optional(),
maxJobsPerMonthLimit: z.number().int().min(1).max(100000).nullable().optional(),
}); });
export const updateRoleSchema = z.object({ export const updateRoleSchema = z.object({
...@@ -30,6 +35,16 @@ export const updateRoleSchema = z.object({ ...@@ -30,6 +35,16 @@ export const updateRoleSchema = z.object({
.optional(), .optional(),
description: z.string().max(500).optional(), description: z.string().max(500).optional(),
isActive: z.boolean().optional(), isActive: z.boolean().optional(),
maxPagesLimit: z.number().int().min(1).max(100000).optional(),
maxJobsPerDayLimit: z.number().int().min(1).max(10000).optional(),
maxConcurrentJobsLimit: z.number().int().min(1).max(100).optional(),
maxPagesPerMonthLimit: z.number().int().min(1).max(1000000).nullable().optional(),
maxJobsPerMonthLimit: z.number().int().min(1).max(100000).nullable().optional(),
syncUsersQuota: z.boolean().optional(),
});
export const resetRoleQuotaSchema = z.object({
syncLimits: z.boolean().optional().default(false),
}); });
export const assignRolePermissionsSchema = z.object({ export const assignRolePermissionsSchema = z.object({
......
...@@ -28,4 +28,39 @@ describe("UserService admin role guard", () => { ...@@ -28,4 +28,39 @@ describe("UserService admin role guard", () => {
); );
expect(repository.update).not.toHaveBeenCalled(); expect(repository.update).not.toHaveBeenCalled();
}); });
it("resets user quota successfully and logs audit event", async () => {
const service = new UserService();
const repository = {
findById: jest.fn().mockResolvedValue(admin),
resetQuota: jest.fn().mockResolvedValue({
...admin,
quotaResetAt: new Date(),
}),
};
const auditLogService = {
log: jest.fn().mockResolvedValue(undefined),
};
(service as any).repository = repository;
(service as any).auditLogService = auditLogService;
const result = await service.resetQuota(admin.id, true, {
actorId: "superadmin-1",
ipAddress: "127.0.0.1",
userAgent: "jest",
});
expect(repository.resetQuota).toHaveBeenCalledWith(admin.id, true);
expect(auditLogService.log).toHaveBeenCalledWith(
expect.objectContaining({
userId: "superadmin-1",
action: "USER_QUOTA_RESET",
details: expect.objectContaining({
targetUserId: admin.id,
resetLimitsToRole: true,
}),
}),
);
expect(result.id).toBe(admin.id);
});
}); });
...@@ -91,6 +91,33 @@ export class UserController { ...@@ -91,6 +91,33 @@ export class UserController {
} }
}; };
resetQuota = async (
req: Request,
res: Response,
next: NextFunction,
): Promise<void> => {
try {
const resetLimitsToRole = !!req.body?.resetLimitsToRole;
const result = await this.service.resetQuota(
req.params.id,
resetLimitsToRole,
{
actorId: req.user.id,
ipAddress: req.ip,
userAgent: req.headers["user-agent"] as string,
},
);
res.json({
success: true,
message: "User quota reset successfully",
data: result,
});
} catch (error) {
next(error);
}
};
delete = async (req: Request, res: Response, next: NextFunction) => { delete = async (req: Request, res: Response, next: NextFunction) => {
try { try {
await this.service.delete( await this.service.delete(
......
...@@ -9,6 +9,8 @@ export interface CreateUserDto { ...@@ -9,6 +9,8 @@ export interface CreateUserDto {
maxPagesLimit?: number; maxPagesLimit?: number;
maxJobsPerDayLimit?: number; maxJobsPerDayLimit?: number;
maxConcurrentJobsLimit?: number; maxConcurrentJobsLimit?: number;
maxPagesPerMonthLimit?: number | null;
maxJobsPerMonthLimit?: number | null;
} }
export interface UpdateUserDto { export interface UpdateUserDto {
...@@ -20,6 +22,12 @@ export interface UpdateUserDto { ...@@ -20,6 +22,12 @@ export interface UpdateUserDto {
maxPagesLimit?: number; maxPagesLimit?: number;
maxJobsPerDayLimit?: number; maxJobsPerDayLimit?: number;
maxConcurrentJobsLimit?: number; maxConcurrentJobsLimit?: number;
maxPagesPerMonthLimit?: number | null;
maxJobsPerMonthLimit?: number | null;
}
export interface ResetUserQuotaDto {
resetLimitsToRole?: boolean;
} }
export interface UserAssignedRoleSummaryDto { export interface UserAssignedRoleSummaryDto {
...@@ -41,6 +49,9 @@ export interface UserResponseDto { ...@@ -41,6 +49,9 @@ export interface UserResponseDto {
maxPagesLimit: number; maxPagesLimit: number;
maxJobsPerDayLimit: number; maxJobsPerDayLimit: number;
maxConcurrentJobsLimit: number; maxConcurrentJobsLimit: number;
maxPagesPerMonthLimit?: number | null;
maxJobsPerMonthLimit?: number | null;
quotaResetAt?: Date | null;
createdAt: Date; createdAt: Date;
updatedAt: Date; updatedAt: Date;
roles?: UserAssignedRoleSummaryDto[]; roles?: UserAssignedRoleSummaryDto[];
......
...@@ -5,6 +5,8 @@ import { envConfig } from "../../config/env.config"; ...@@ -5,6 +5,8 @@ import { envConfig } from "../../config/env.config";
import { ROLES } from "../../common/constants/role.constant"; import { ROLES } from "../../common/constants/role.constant";
import { SYSTEM_ROLE_SLUGS } from "../../common/constants/system-role.constant"; import { SYSTEM_ROLE_SLUGS } from "../../common/constants/system-role.constant";
import { systemConfigService } from "../system-config/system-config.service"; import { systemConfigService } from "../system-config/system-config.service";
import { AppError } from "../../common/errors/app-error";
import { ERROR_CODE } from "../../common/errors/error-code";
export class UserRepository { export class UserRepository {
async findAll(query: UserQueryDto = {}) { async findAll(query: UserQueryDto = {}) {
...@@ -104,6 +106,8 @@ export class UserRepository { ...@@ -104,6 +106,8 @@ export class UserRepository {
maxPagesLimit?: number; maxPagesLimit?: number;
maxJobsPerDayLimit?: number; maxJobsPerDayLimit?: number;
maxConcurrentJobsLimit?: number; maxConcurrentJobsLimit?: number;
maxPagesPerMonthLimit?: number | null;
maxJobsPerMonthLimit?: number | null;
}): Promise<User> { }): Promise<User> {
const defaultMaxPages = await systemConfigService.get<number>( const defaultMaxPages = await systemConfigService.get<number>(
"quota.user_max_pages", "quota.user_max_pages",
...@@ -117,6 +121,14 @@ export class UserRepository { ...@@ -117,6 +121,14 @@ export class UserRepository {
"quota.user_max_concurrent_jobs", "quota.user_max_concurrent_jobs",
envConfig.quota.defaultMaxConcurrentJobs, envConfig.quota.defaultMaxConcurrentJobs,
); );
const defaultMaxPagesPerMonth = await systemConfigService.get<number>(
"quota.user_max_pages_per_month",
envConfig.quota.defaultMaxPagesPerMonth,
);
const defaultMaxJobsPerMonth = await systemConfigService.get<number>(
"quota.user_max_jobs_per_month",
envConfig.quota.defaultMaxJobsPerMonth,
);
return prisma.user.create({ return prisma.user.create({
data: { data: {
...@@ -129,6 +141,10 @@ export class UserRepository { ...@@ -129,6 +141,10 @@ export class UserRepository {
maxJobsPerDayLimit: data.maxJobsPerDayLimit ?? defaultMaxJobsPerDay, maxJobsPerDayLimit: data.maxJobsPerDayLimit ?? defaultMaxJobsPerDay,
maxConcurrentJobsLimit: maxConcurrentJobsLimit:
data.maxConcurrentJobsLimit ?? defaultMaxConcurrentJobs, data.maxConcurrentJobsLimit ?? defaultMaxConcurrentJobs,
maxPagesPerMonthLimit:
data.maxPagesPerMonthLimit ?? defaultMaxPagesPerMonth,
maxJobsPerMonthLimit:
data.maxJobsPerMonthLimit ?? defaultMaxJobsPerMonth,
}, },
}); });
} }
...@@ -144,6 +160,9 @@ export class UserRepository { ...@@ -144,6 +160,9 @@ export class UserRepository {
maxPagesLimit?: number; maxPagesLimit?: number;
maxJobsPerDayLimit?: number; maxJobsPerDayLimit?: number;
maxConcurrentJobsLimit?: number; maxConcurrentJobsLimit?: number;
maxPagesPerMonthLimit?: number | null;
maxJobsPerMonthLimit?: number | null;
quotaResetAt?: Date | null;
}, },
): Promise<User> { ): Promise<User> {
return prisma.user.update({ return prisma.user.update({
...@@ -152,6 +171,65 @@ export class UserRepository { ...@@ -152,6 +171,65 @@ export class UserRepository {
}); });
} }
async resetQuota(
id: string,
resetLimitsToRole: boolean = false,
): Promise<User> {
const now = new Date();
if (!resetLimitsToRole) {
return prisma.user.update({
where: { id },
data: { quotaResetAt: now },
include: {
userRoles: {
include: { role: true },
},
},
});
}
const user = await prisma.user.findUnique({
where: { id },
include: {
userRoles: {
include: { role: true },
},
},
});
if (!user) {
throw new AppError("User not found", 404, ERROR_CODE.NOT_FOUND);
}
const primaryRole =
user.userRoles.find((ur) => ur.role.isActive)?.role ||
(await prisma.role.findFirst({
where: { slug: user.role.toLowerCase() },
}));
const updateData: any = {
quotaResetAt: now,
};
if (primaryRole) {
updateData.maxPagesLimit = primaryRole.maxPagesLimit;
updateData.maxJobsPerDayLimit = primaryRole.maxJobsPerDayLimit;
updateData.maxConcurrentJobsLimit = primaryRole.maxConcurrentJobsLimit;
updateData.maxPagesPerMonthLimit = primaryRole.maxPagesPerMonthLimit;
updateData.maxJobsPerMonthLimit = primaryRole.maxJobsPerMonthLimit;
}
return prisma.user.update({
where: { id },
data: updateData,
include: {
userRoles: {
include: { role: true },
},
},
});
}
async delete(id: string, deletedBy: string): Promise<User> { async delete(id: string, deletedBy: string): Promise<User> {
return prisma.$transaction(async (tx) => { return prisma.$transaction(async (tx) => {
const user = await tx.user.update({ const user = await tx.user.update({
......
...@@ -10,6 +10,7 @@ import { ...@@ -10,6 +10,7 @@ import {
import { import {
createUserSchema, createUserSchema,
updateUserSchema, updateUserSchema,
resetUserQuotaSchema,
listUsersQuerySchema, listUsersQuerySchema,
assignUserRolesSchema, assignUserRolesSchema,
userParamsSchema, userParamsSchema,
...@@ -67,6 +68,15 @@ router.delete( ...@@ -67,6 +68,15 @@ router.delete(
controller.delete, controller.delete,
); );
router.post(
"/:id/reset-quota",
authMiddleware,
requirePermission(PERMISSIONS.USERS_UPDATE),
validateParams(userParamsSchema),
validate(resetUserQuotaSchema),
controller.resetQuota,
);
router.get( router.get(
"/:id/roles", "/:id/roles",
authMiddleware, authMiddleware,
......
...@@ -51,6 +51,9 @@ export class UserService { ...@@ -51,6 +51,9 @@ export class UserService {
maxPagesLimit: user.maxPagesLimit, maxPagesLimit: user.maxPagesLimit,
maxJobsPerDayLimit: user.maxJobsPerDayLimit, maxJobsPerDayLimit: user.maxJobsPerDayLimit,
maxConcurrentJobsLimit: user.maxConcurrentJobsLimit, maxConcurrentJobsLimit: user.maxConcurrentJobsLimit,
maxPagesPerMonthLimit: user.maxPagesPerMonthLimit ?? null,
maxJobsPerMonthLimit: user.maxJobsPerMonthLimit ?? null,
quotaResetAt: user.quotaResetAt ?? null,
createdAt: user.createdAt, createdAt: user.createdAt,
updatedAt: user.updatedAt, updatedAt: user.updatedAt,
...(roles !== undefined ? { roles } : {}), ...(roles !== undefined ? { roles } : {}),
...@@ -101,6 +104,8 @@ export class UserService { ...@@ -101,6 +104,8 @@ export class UserService {
maxPagesLimit: data.maxPagesLimit, maxPagesLimit: data.maxPagesLimit,
maxJobsPerDayLimit: data.maxJobsPerDayLimit, maxJobsPerDayLimit: data.maxJobsPerDayLimit,
maxConcurrentJobsLimit: data.maxConcurrentJobsLimit, maxConcurrentJobsLimit: data.maxConcurrentJobsLimit,
maxPagesPerMonthLimit: data.maxPagesPerMonthLimit,
maxJobsPerMonthLimit: data.maxJobsPerMonthLimit,
}); });
// Auto assign matching default system role // Auto assign matching default system role
...@@ -190,12 +195,43 @@ export class UserService { ...@@ -190,12 +195,43 @@ export class UserService {
maxPagesLimit: data.maxPagesLimit, maxPagesLimit: data.maxPagesLimit,
maxJobsPerDayLimit: data.maxJobsPerDayLimit, maxJobsPerDayLimit: data.maxJobsPerDayLimit,
maxConcurrentJobsLimit: data.maxConcurrentJobsLimit, maxConcurrentJobsLimit: data.maxConcurrentJobsLimit,
maxPagesPerMonthLimit: data.maxPagesPerMonthLimit,
maxJobsPerMonthLimit: data.maxJobsPerMonthLimit,
}); });
authorizationCache.invalidateUser(id); authorizationCache.invalidateUser(id);
return this.formatUser(user); return this.formatUser(user);
} }
async resetQuota(
id: string,
resetLimitsToRole?: boolean,
context?: { actorId?: string; ipAddress?: string; userAgent?: string },
): Promise<UserResponseDto> {
const user = await this.repository.findById(id);
if (!user) {
throw new AppError("User not found", 404, ERROR_CODE.NOT_FOUND);
}
const updated = await this.repository.resetQuota(id, resetLimitsToRole);
if (context?.actorId) {
await this.auditLogService.log({
userId: context.actorId,
action: AUDIT_ACTIONS.USER_QUOTA_RESET,
details: {
targetUserId: id,
targetUserEmail: user.email,
resetLimitsToRole: !!resetLimitsToRole,
},
ipAddress: context.ipAddress,
userAgent: context.userAgent,
});
}
return this.formatUser(updated);
}
async delete( async delete(
id: string, id: string,
currentUserId: string, currentUserId: string,
......
...@@ -14,6 +14,8 @@ export const createUserSchema = z.object({ ...@@ -14,6 +14,8 @@ export const createUserSchema = z.object({
maxPagesLimit: z.number().int().min(1).max(100000).optional(), maxPagesLimit: z.number().int().min(1).max(100000).optional(),
maxJobsPerDayLimit: z.number().int().min(1).max(10000).optional(), maxJobsPerDayLimit: z.number().int().min(1).max(10000).optional(),
maxConcurrentJobsLimit: z.number().int().min(1).max(100).optional(), maxConcurrentJobsLimit: z.number().int().min(1).max(100).optional(),
maxPagesPerMonthLimit: z.number().int().min(1).max(1000000).nullable().optional(),
maxJobsPerMonthLimit: z.number().int().min(1).max(100000).nullable().optional(),
}); });
export const updateUserSchema = z.object({ export const updateUserSchema = z.object({
...@@ -30,6 +32,12 @@ export const updateUserSchema = z.object({ ...@@ -30,6 +32,12 @@ export const updateUserSchema = z.object({
maxPagesLimit: z.number().int().min(1).max(100000).optional(), maxPagesLimit: z.number().int().min(1).max(100000).optional(),
maxJobsPerDayLimit: z.number().int().min(1).max(10000).optional(), maxJobsPerDayLimit: z.number().int().min(1).max(10000).optional(),
maxConcurrentJobsLimit: z.number().int().min(1).max(100).optional(), maxConcurrentJobsLimit: z.number().int().min(1).max(100).optional(),
maxPagesPerMonthLimit: z.number().int().min(1).max(1000000).nullable().optional(),
maxJobsPerMonthLimit: z.number().int().min(1).max(100000).nullable().optional(),
});
export const resetUserQuotaSchema = z.object({
resetLimitsToRole: z.boolean().optional().default(false),
}); });
export const listUsersQuerySchema = z.object({ export const listUsersQuerySchema = z.object({
......
import { Queue } from "bullmq";
import { envConfig } from "../config/env.config";
import {
CRON_QUEUE_NAME,
CronJobName,
DEFAULT_CRON_SCHEDULES,
DEFAULT_CRON_TIMEZONE,
} from "../common/constants/cron.constant";
export class CronQueueService {
private queue: Queue | null = null;
constructor() {
if (envConfig.redis.enabled) {
this.queue = new Queue(CRON_QUEUE_NAME, {
connection: {
host: envConfig.redis.host,
port: envConfig.redis.port,
enableOfflineQueue: false,
lazyConnect: true,
},
defaultJobOptions: {
attempts: 3,
backoff: { type: "exponential", delay: 5000 },
removeOnComplete: { count: 200 },
removeOnFail: { count: 1000 },
},
});
}
}
getQueue(): Queue | null {
return this.queue;
}
/**
* Đăng ký hoặc đồng bộ toàn bộ lịch biểu mặc định với BullMQ Schedulers
*/
async registerDefaultSchedulers(
jobStatuses: Record<string, boolean>,
): Promise<void> {
if (!this.queue) return;
for (const name of Object.keys(DEFAULT_CRON_SCHEDULES)) {
const jobName = name as CronJobName;
const isEnabled = jobStatuses[jobName] ?? true;
try {
if (isEnabled) {
await this.enableJobScheduler(jobName);
} else {
await this.disableJobScheduler(jobName);
}
} catch (err) {
console.warn(
`[Cron Queue] Failed to register schedule for ${jobName}:`,
err,
);
}
}
}
/**
* Bật lịch trình chạy định kỳ tự động trong BullMQ Scheduler
*/
async enableJobScheduler(jobName: CronJobName): Promise<void> {
if (!this.queue) return;
const config = DEFAULT_CRON_SCHEDULES[jobName];
if (!config) return;
try {
await this.queue.upsertJobScheduler(
jobName,
{
pattern: config.cron,
tz: DEFAULT_CRON_TIMEZONE,
},
{
name: jobName,
data: { params: config.defaultParams },
},
);
} catch (err) {
console.warn(
`[Cron Queue] Could not upsert scheduler for ${jobName}:`,
err,
);
}
}
/**
* Tắt lịch trình chạy định kỳ tự động khỏi BullMQ Scheduler
*/
async disableJobScheduler(jobName: CronJobName): Promise<void> {
if (!this.queue) return;
try {
await this.queue.removeJobScheduler(jobName);
} catch (err) {
console.warn(
`[Cron Queue] Could not remove scheduler for ${jobName}:`,
err,
);
}
}
}
export const cronQueue = new CronQueueService();
import "dotenv/config";
import { Worker } from "bullmq";
import { envConfig } from "../config/env.config";
import { cronService } from "../modules/cron/cron.service";
import {
CRON_QUEUE_NAME,
CronJobName,
} from "../common/constants/cron.constant";
import { getErrorMessage } from "../common/helpers/error-mapping.helper";
if (!envConfig.redis.enabled) {
console.log(
"[Cron Worker] REDIS_ENABLED is not set to true. Cron Worker will not start.",
);
process.exit(0);
}
export const cronWorker = new Worker(
CRON_QUEUE_NAME,
async (job) => {
const jobName = job.name as CronJobName;
console.log(`[Cron Worker] Processing periodic job: ${jobName} (ID: ${job.id})`);
// 1. Kiểm tra trạng thái bật/tắt của tác vụ từ SystemConfig
const isEnabled = await cronService.isJobEnabled(jobName);
if (!isEnabled) {
console.log(
`[Cron Worker] Job '${jobName}' is currently DISABLED in system configuration. Skipping execution.`,
);
return { skipped: true, reason: "JOB_DISABLED" };
}
// 2. Thực thi nghiệp vụ
try {
const result = await cronService.executeJob(
jobName,
job.data?.params as Record<string, unknown> | undefined,
{ source: "SCHEDULER" },
);
if (!result.success) {
throw new Error(result.error || `Tác vụ '${jobName}' thực thi thất bại`);
}
console.log(
`[Cron Worker] Job '${jobName}' completed successfully in ${result.durationMs}ms`,
);
return result;
} catch (err: unknown) {
const errorMessage = getErrorMessage(err);
console.error(
`[Cron Worker] Job '${jobName}' execution failed: ${errorMessage}`,
);
throw err;
}
},
{
connection: {
host: envConfig.redis.host,
port: envConfig.redis.port,
maxRetriesPerRequest: null,
},
concurrency: 2,
},
);
cronWorker.on("error", (err) => {
console.error("[Cron Worker] Error:", err);
});
export async function closeCronWorker() {
console.log("[Cron Worker] Closing worker gracefully...");
await cronWorker.close();
console.log("[Cron Worker] Closed");
}
// Standalone execution
if (require.main === module) {
const gracefulShutdown = async (signal: string) => {
console.log(
`[Cron Worker] Received ${signal}, initiating graceful shutdown...`,
);
await closeCronWorker();
process.exit(0);
};
process.on("SIGTERM", () => gracefulShutdown("SIGTERM"));
process.on("SIGINT", () => gracefulShutdown("SIGINT"));
console.log("[Cron Worker] Standalone Cron Worker started");
}
...@@ -13,6 +13,7 @@ import dashboardRoute from "../modules/dashboard/dashboard.route"; ...@@ -13,6 +13,7 @@ import dashboardRoute from "../modules/dashboard/dashboard.route";
import roleRoute from "../modules/roles/role.route"; import roleRoute from "../modules/roles/role.route";
import permissionRoute from "../modules/permissions/permission.route"; import permissionRoute from "../modules/permissions/permission.route";
import systemConfigRoute from "../modules/system-config/system-config.route"; import systemConfigRoute from "../modules/system-config/system-config.route";
import cronRoute from "../modules/cron/cron.route";
const router = Router(); const router = Router();
...@@ -22,6 +23,7 @@ router.use("/users", userRoute); ...@@ -22,6 +23,7 @@ router.use("/users", userRoute);
router.use("/roles", roleRoute); router.use("/roles", roleRoute);
router.use("/permissions", permissionRoute); router.use("/permissions", permissionRoute);
router.use("/system", systemConfigRoute); router.use("/system", systemConfigRoute);
router.use("/cron", cronRoute);
router.use("/dashboard", dashboardRoute); router.use("/dashboard", dashboardRoute);
router.use("/crawl-jobs", crawlJobRoute); router.use("/crawl-jobs", crawlJobRoute);
router.use("/crawl-schedules", crawlScheduleRoute); router.use("/crawl-schedules", crawlScheduleRoute);
...@@ -31,3 +33,4 @@ router.use("/api-keys", apiKeyRoute); ...@@ -31,3 +33,4 @@ router.use("/api-keys", apiKeyRoute);
router.use("/webhooks", webhookRoute); router.use("/webhooks", webhookRoute);
router.use("/extraction-templates", extractionTemplateRoute); router.use("/extraction-templates", extractionTemplateRoute);
export default router; export default router;
...@@ -55,6 +55,16 @@ async function bootstrap() { ...@@ -55,6 +55,16 @@ async function bootstrap() {
console.warn("[Server] Failed to initialize default system configs:", err); console.warn("[Server] Failed to initialize default system configs:", err);
} }
const { permissionService } = await import(
"./modules/permissions/permission.service"
);
try {
await permissionService.ensureSystemPermissions();
console.log("[Server] System permissions and role bindings synchronized successfully.");
} catch (err) {
console.warn("[Server] Failed to synchronize system permissions:", err);
}
if (isRedisAvailable) { if (isRedisAvailable) {
systemConfigService.initRedisSubscriber(); systemConfigService.initRedisSubscriber();
...@@ -64,6 +74,19 @@ async function bootstrap() { ...@@ -64,6 +74,19 @@ async function bootstrap() {
const { startScheduleWorker } = await import("./queues/schedule.worker"); const { startScheduleWorker } = await import("./queues/schedule.worker");
startScheduleWorker(); startScheduleWorker();
console.log("[Server] Schedule worker initialized in background."); console.log("[Server] Schedule worker initialized in background.");
const { cronQueue } = await import("./queues/cron.queue");
const { cronRepository } = await import("./modules/cron/cron.repository");
try {
const jobStatuses = await cronRepository.getJobStatuses();
await cronQueue.registerDefaultSchedulers(jobStatuses);
console.log("[Server] Cron job schedulers registered successfully.");
} catch (cronErr) {
console.warn("[Server] Failed to register cron schedulers:", cronErr);
}
await import("./queues/cron.worker");
console.log("[Server] Cron worker initialized in background.");
} }
app.listen(envConfig.port, () => { app.listen(envConfig.port, () => {
......
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment