From 3ed5dd61093833a6dcdd314a90402f6ec292c918 Mon Sep 17 00:00:00 2001 From: = Date: Thu, 16 May 2024 15:41:03 +0530 Subject: [PATCH] feat: removed audit log queue and switched to resource clean up queue --- .../ee/services/audit-log/audit-log-queue.ts | 31 +--------- backend/src/queue/queue-service.ts | 12 +++- backend/src/server/routes/index.ts | 8 ++- .../resource-cleanup-queue.ts | 58 +++++++++++++++++++ 4 files changed, 77 insertions(+), 32 deletions(-) create mode 100644 backend/src/services/resource-cleanup/resource-cleanup-queue.ts diff --git a/backend/src/ee/services/audit-log/audit-log-queue.ts b/backend/src/ee/services/audit-log/audit-log-queue.ts index 6c563b573..f93b391a5 100644 --- a/backend/src/ee/services/audit-log/audit-log-queue.ts +++ b/backend/src/ee/services/audit-log/audit-log-queue.ts @@ -3,7 +3,6 @@ import { RawAxiosRequestHeaders } from "axios"; import { SecretKeyEncoding } from "@app/db/schemas"; import { request } from "@app/lib/config/request"; import { infisicalSymmetricDecrypt } from "@app/lib/crypto/encryption"; -import { logger } from "@app/lib/logger"; import { QueueJobs, QueueName, TQueueServiceFactory } from "@app/queue"; import { TProjectDALFactory } from "@app/services/project/project-dal"; @@ -113,35 +112,7 @@ export const auditLogQueueServiceFactory = ({ ); }); - queueService.start(QueueName.AuditLogPrune, async () => { - logger.info(`${QueueName.AuditLogPrune}: queue task started`); - await auditLogDAL.pruneAuditLog(); - logger.info(`${QueueName.AuditLogPrune}: queue task completed`); - }); - - // we do a repeat cron job in utc timezone at 12 Midnight each day - const startAuditLogPruneJob = async () => { - // clear previous job - await queueService.stopRepeatableJob( - QueueName.AuditLogPrune, - QueueJobs.AuditLogPrune, - { pattern: "0 0 * * *", utc: true }, - QueueName.AuditLogPrune // just a job id - ); - - await queueService.queue(QueueName.AuditLogPrune, QueueJobs.AuditLogPrune, undefined, { - delay: 5000, - jobId: QueueName.AuditLogPrune, - repeat: { pattern: "0 0 * * *", utc: true } - }); - }; - - queueService.listen(QueueName.AuditLogPrune, "failed", (err) => { - logger.error(err?.failedReason, `${QueueName.AuditLogPrune}: log pruning failed`); - }); - return { - pushToLog, - startAuditLogPruneJob + pushToLog }; }; diff --git a/backend/src/queue/queue-service.ts b/backend/src/queue/queue-service.ts index bc8ac88ff..9d85b6015 100644 --- a/backend/src/queue/queue-service.ts +++ b/backend/src/queue/queue-service.ts @@ -12,7 +12,9 @@ export enum QueueName { SecretRotation = "secret-rotation", SecretReminder = "secret-reminder", AuditLog = "audit-log", + // TODO(akhilmhdh): This will get removed later. For now this is kept to stop the repeatable queue AuditLogPrune = "audit-log-prune", + DailyResourceCleanUp = "daily-resource-cleanup", TelemetryInstanceStats = "telemtry-self-hosted-stats", IntegrationSync = "sync-integrations", SecretWebhook = "secret-webhook", @@ -26,7 +28,9 @@ export enum QueueJobs { SecretReminder = "secret-reminder-job", SecretRotation = "secret-rotation-job", AuditLog = "audit-log-job", + // TODO(akhilmhdh): This will get removed later. For now this is kept to stop the repeatable queue AuditLogPrune = "audit-log-prune-job", + DailyResourceCleanUp = "daily-resource-cleanup-job", SecWebhook = "secret-webhook-trigger", TelemetryInstanceStats = "telemetry-self-hosted-stats", IntegrationSync = "secret-integration-pull", @@ -55,6 +59,10 @@ export type TQueueJobTypes = { name: QueueJobs.AuditLog; payload: TCreateAuditLogDTO; }; + [QueueName.DailyResourceCleanUp]: { + name: QueueJobs.DailyResourceCleanUp; + payload: undefined; + }; [QueueName.AuditLogPrune]: { name: QueueJobs.AuditLogPrune; payload: undefined; @@ -172,7 +180,9 @@ export const queueServiceFactory = (redisUrl: string) => { jobId?: string ) => { const q = queueContainer[name]; - return q.removeRepeatable(job, repeatOpt, jobId); + if (q) { + return q.removeRepeatable(job, repeatOpt, jobId); + } }; const stopRepeatableJobByJobId = async (name: T, jobId: string) => { diff --git a/backend/src/server/routes/index.ts b/backend/src/server/routes/index.ts index 3a050bcad..8ef6ee3be 100644 --- a/backend/src/server/routes/index.ts +++ b/backend/src/server/routes/index.ts @@ -113,6 +113,7 @@ import { projectMembershipServiceFactory } from "@app/services/project-membershi import { projectUserMembershipRoleDALFactory } from "@app/services/project-membership/project-user-membership-role-dal"; import { projectRoleDALFactory } from "@app/services/project-role/project-role-dal"; import { projectRoleServiceFactory } from "@app/services/project-role/project-role-service"; +import { dailyResourceCleanUpQueueServiceFactory } from "@app/services/resource-cleanup/resource-cleanup-queue"; import { secretDALFactory } from "@app/services/secret/secret-dal"; import { secretQueueFactory } from "@app/services/secret/secret-queue"; import { secretServiceFactory } from "@app/services/secret/secret-service"; @@ -757,14 +758,19 @@ export const registerRoutes = async ( folderDAL, licenseService }); + const dailyResourceCleanUp = dailyResourceCleanUpQueueServiceFactory({ + auditLogDAL, + queueService, + identityAccessTokenDAL + }); await superAdminService.initServerCfg(); // // setup the communication with license key server await licenseService.init(); - await auditLogQueue.startAuditLogPruneJob(); await telemetryQueue.startTelemetryCheck(); + await dailyResourceCleanUp.startCleanUp(); // inject all services server.decorate("services", { diff --git a/backend/src/services/resource-cleanup/resource-cleanup-queue.ts b/backend/src/services/resource-cleanup/resource-cleanup-queue.ts new file mode 100644 index 000000000..3c8bcb1f7 --- /dev/null +++ b/backend/src/services/resource-cleanup/resource-cleanup-queue.ts @@ -0,0 +1,58 @@ +import { TAuditLogDALFactory } from "@app/ee/services/audit-log/audit-log-dal"; +import { logger } from "@app/lib/logger"; +import { QueueJobs, QueueName, TQueueServiceFactory } from "@app/queue"; + +import { TIdentityAccessTokenDALFactory } from "../identity-access-token/identity-access-token-dal"; + +type TDailyResourceCleanUpQueueServiceFactoryDep = { + auditLogDAL: Pick; + identityAccessTokenDAL: Pick; + queueService: TQueueServiceFactory; +}; + +export type TDailyResourceCleanUpQueueServiceFactory = ReturnType; + +export const dailyResourceCleanUpQueueServiceFactory = ({ + auditLogDAL, + queueService, + identityAccessTokenDAL +}: TDailyResourceCleanUpQueueServiceFactoryDep) => { + queueService.start(QueueName.DailyResourceCleanUp, async () => { + logger.info(`${QueueName.DailyResourceCleanUp}: queue task started`); + await auditLogDAL.pruneAuditLog(); + await identityAccessTokenDAL.removeExpiredTokens(); + logger.info(`${QueueName.DailyResourceCleanUp}: queue task completed`); + }); + + // we do a repeat cron job in utc timezone at 12 Midnight each day + const startCleanUp = async () => { + // TODO(akhilmhdh): remove later + await queueService.stopRepeatableJob( + QueueName.AuditLogPrune, + QueueJobs.AuditLogPrune, + { pattern: "0 0 * * *", utc: true }, + QueueName.AuditLogPrune // just a job id + ); + // clear previous job + await queueService.stopRepeatableJob( + QueueName.DailyResourceCleanUp, + QueueJobs.DailyResourceCleanUp, + { pattern: "0 0 * * *", utc: true }, + QueueName.DailyResourceCleanUp // just a job id + ); + + await queueService.queue(QueueName.DailyResourceCleanUp, QueueJobs.DailyResourceCleanUp, undefined, { + delay: 5000, + jobId: QueueName.DailyResourceCleanUp, + repeat: { pattern: "0 0 * * *", utc: true } + }); + }; + + queueService.listen(QueueName.DailyResourceCleanUp, "failed", (_, err) => { + logger.error(err, `${QueueName.DailyResourceCleanUp}: resource cleanup failed`); + }); + + return { + startCleanUp + }; +};