diff --git a/backend/src/lib/config/env.ts b/backend/src/lib/config/env.ts index 94b1e4a44..d44eb90d6 100644 --- a/backend/src/lib/config/env.ts +++ b/backend/src/lib/config/env.ts @@ -79,6 +79,7 @@ const envSchema = z QUEUE_WORKER_PROFILE: z.nativeEnum(QueueWorkerProfile).default(QueueWorkerProfile.All), HTTPS_ENABLED: zodStrBool, ROTATION_DEVELOPMENT_MODE: zodStrBool.default("false").optional(), + DAILY_RESOURCE_CLEAN_UP_DEVELOPMENT_MODE: zodStrBool.default("false").optional(), // smtp options SMTP_HOST: zpStr(z.string().optional()), SMTP_IGNORE_TLS: zodStrBool.default("false"), @@ -348,6 +349,8 @@ const envSchema = z isRedisConfigured: Boolean(data.REDIS_URL || data.REDIS_SENTINEL_HOSTS), isDevelopmentMode: data.NODE_ENV === "development", isRotationDevelopmentMode: data.NODE_ENV === "development" && data.ROTATION_DEVELOPMENT_MODE, + isDailyResourceCleanUpDevelopmentMode: + data.NODE_ENV === "development" && data.DAILY_RESOURCE_CLEAN_UP_DEVELOPMENT_MODE, isProductionMode: data.NODE_ENV === "production" || IS_PACKAGED, isRedisSentinelMode: Boolean(data.REDIS_SENTINEL_HOSTS), REDIS_SENTINEL_HOSTS: data.REDIS_SENTINEL_HOSTS?.trim() diff --git a/backend/src/server/routes/index.ts b/backend/src/server/routes/index.ts index 90e961c29..8be9f3034 100644 --- a/backend/src/server/routes/index.ts +++ b/backend/src/server/routes/index.ts @@ -1972,7 +1972,7 @@ export const registerRoutes = async ( await telemetryQueue.startTelemetryCheck(); await telemetryQueue.startAggregatedEventsJob(); - await dailyResourceCleanUp.startCleanUp(); + await dailyResourceCleanUp.init(); await dailyReminderQueueService.startDailyRemindersJob(); await dailyReminderQueueService.startSecretReminderMigrationJob(); await dailyExpiringPkiItemAlert.startSendingAlerts(); diff --git a/backend/src/services/resource-cleanup/resource-cleanup-queue.ts b/backend/src/services/resource-cleanup/resource-cleanup-queue.ts index bfaef5708..224eef4bf 100644 --- a/backend/src/services/resource-cleanup/resource-cleanup-queue.ts +++ b/backend/src/services/resource-cleanup/resource-cleanup-queue.ts @@ -1,5 +1,6 @@ import { TAuditLogDALFactory } from "@app/ee/services/audit-log/audit-log-dal"; import { TSnapshotDALFactory } from "@app/ee/services/secret-snapshot/snapshot-dal"; +import { getConfig } from "@app/lib/config/env"; import { logger } from "@app/lib/logger"; import { QueueJobs, QueueName, TQueueServiceFactory } from "@app/queue"; @@ -41,51 +42,50 @@ export const dailyResourceCleanUpQueueServiceFactory = ({ serviceTokenService, orgService }: TDailyResourceCleanUpQueueServiceFactoryDep) => { - queueService.start(QueueName.DailyResourceCleanUp, async () => { - logger.info(`${QueueName.DailyResourceCleanUp}: queue task started`); - await identityAccessTokenDAL.removeExpiredTokens(); - await identityUniversalAuthClientSecretDAL.removeExpiredClientSecrets(); - await secretSharingDAL.pruneExpiredSharedSecrets(); - await secretSharingDAL.pruneExpiredSecretRequests(); - await snapshotDAL.pruneExcessSnapshots(); - await secretVersionDAL.pruneExcessVersions(); - await secretVersionV2DAL.pruneExcessVersions(); - await secretFolderVersionDAL.pruneExcessVersions(); - await serviceTokenService.notifyExpiringTokens(); - await orgService.notifyInvitedUsers(); - await auditLogDAL.pruneAuditLog(); - logger.info(`${QueueName.DailyResourceCleanUp}: queue task completed`); - }); + const appCfg = getConfig(); - // 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, + if (appCfg.isDailyResourceCleanUpDevelopmentMode) { + logger.warn("Daily Resource Clean Up is in development mode."); + } + + const init = async () => { + await queueService.startPg( QueueJobs.DailyResourceCleanUp, - { pattern: "0 0 * * *", utc: true }, - QueueName.DailyResourceCleanUp // just a job id + async () => { + try { + logger.info(`${QueueName.DailyResourceCleanUp}: queue task started`); + await identityAccessTokenDAL.removeExpiredTokens(); + await identityUniversalAuthClientSecretDAL.removeExpiredClientSecrets(); + await secretSharingDAL.pruneExpiredSharedSecrets(); + await secretSharingDAL.pruneExpiredSecretRequests(); + await snapshotDAL.pruneExcessSnapshots(); + await secretVersionDAL.pruneExcessVersions(); + await secretVersionV2DAL.pruneExcessVersions(); + await secretFolderVersionDAL.pruneExcessVersions(); + await serviceTokenService.notifyExpiringTokens(); + await orgService.notifyInvitedUsers(); + await auditLogDAL.pruneAuditLog(); + logger.info(`${QueueName.DailyResourceCleanUp}: queue task completed`); + } catch (error) { + logger.error(error, `${QueueName.DailyResourceCleanUp}: resource cleanup failed`); + throw error; + } + }, + { + batchSize: 1, + workerCount: 1, + pollingIntervalSeconds: 1 + } + ); + await queueService.schedulePg( + QueueJobs.DailyResourceCleanUp, + appCfg.isDailyResourceCleanUpDevelopmentMode ? "*/5 * * * *" : "0 0 * * *", + undefined, + { tz: "UTC" } ); - - 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 + init }; };