diff --git a/backend/src/queue/queue-service.ts b/backend/src/queue/queue-service.ts index 051fe9cbd..330193052 100644 --- a/backend/src/queue/queue-service.ts +++ b/backend/src/queue/queue-service.ts @@ -317,6 +317,13 @@ export const queueServiceFactory = ( } }; + const getRepeatableJobs = (name: QueueName, startOffset?: number, endOffset?: number) => { + const q = queueContainer[name]; + if (!q) throw new Error(`Queue '${name}' not initialized`); + + return q.getRepeatableJobs(startOffset, endOffset); + }; + const stopRepeatableJobByJobId = async (name: T, jobId: string) => { const q = queueContainer[name]; const job = await q.getJob(jobId); @@ -326,6 +333,11 @@ export const queueServiceFactory = ( return q.removeRepeatableByKey(job.repeatJobKey); }; + const stopRepeatableJobByKey = async (name: T, repeatJobKey: string) => { + const q = queueContainer[name]; + return q.removeRepeatableByKey(repeatJobKey); + }; + const stopJobById = async (name: T, jobId: string) => { const q = queueContainer[name]; const job = await q.getJob(jobId); @@ -349,8 +361,10 @@ export const queueServiceFactory = ( shutdown, stopRepeatableJob, stopRepeatableJobByJobId, + stopRepeatableJobByKey, clearQueue, stopJobById, + getRepeatableJobs, startPg, queuePg }; diff --git a/backend/src/server/routes/index.ts b/backend/src/server/routes/index.ts index dd6520875..65a878bbf 100644 --- a/backend/src/server/routes/index.ts +++ b/backend/src/server/routes/index.ts @@ -1240,6 +1240,7 @@ export const registerRoutes = async ( auditLogDAL, queueService, secretVersionDAL, + secretDAL, secretFolderVersionDAL: folderVersionDAL, snapshotDAL, identityAccessTokenDAL, diff --git a/backend/src/services/resource-cleanup/resource-cleanup-queue.ts b/backend/src/services/resource-cleanup/resource-cleanup-queue.ts index dab70806f..aa1ed9d25 100644 --- a/backend/src/services/resource-cleanup/resource-cleanup-queue.ts +++ b/backend/src/services/resource-cleanup/resource-cleanup-queue.ts @@ -5,6 +5,7 @@ import { QueueJobs, QueueName, TQueueServiceFactory } from "@app/queue"; import { TIdentityAccessTokenDALFactory } from "../identity-access-token/identity-access-token-dal"; import { TIdentityUaClientSecretDALFactory } from "../identity-ua/identity-ua-client-secret-dal"; +import { TSecretDALFactory } from "../secret/secret-dal"; import { TSecretVersionDALFactory } from "../secret/secret-version-dal"; import { TSecretFolderVersionDALFactory } from "../secret-folder/secret-folder-version-dal"; import { TSecretSharingDALFactory } from "../secret-sharing/secret-sharing-dal"; @@ -16,6 +17,7 @@ type TDailyResourceCleanUpQueueServiceFactoryDep = { identityUniversalAuthClientSecretDAL: Pick; secretVersionDAL: Pick; secretVersionV2DAL: Pick; + secretDAL: Pick; secretFolderVersionDAL: Pick; snapshotDAL: Pick; secretSharingDAL: Pick; @@ -30,6 +32,7 @@ export const dailyResourceCleanUpQueueServiceFactory = ({ snapshotDAL, secretVersionDAL, secretFolderVersionDAL, + secretDAL, identityAccessTokenDAL, secretSharingDAL, secretVersionV2DAL, @@ -37,6 +40,7 @@ export const dailyResourceCleanUpQueueServiceFactory = ({ }: TDailyResourceCleanUpQueueServiceFactoryDep) => { queueService.start(QueueName.DailyResourceCleanUp, async () => { logger.info(`${QueueName.DailyResourceCleanUp}: queue task started`); + await secretDAL.pruneSecretReminders(queueService); await auditLogDAL.pruneAuditLog(); await identityAccessTokenDAL.removeExpiredTokens(); await identityUniversalAuthClientSecretDAL.removeExpiredClientSecrets(); diff --git a/backend/src/services/secret/secret-dal.ts b/backend/src/services/secret/secret-dal.ts index 0d4ae0cda..cbaf7ddcd 100644 --- a/backend/src/services/secret/secret-dal.ts +++ b/backend/src/services/secret/secret-dal.ts @@ -5,6 +5,8 @@ import { TDbClient } from "@app/db"; import { SecretsSchema, SecretType, TableName, TSecrets, TSecretsUpdate } from "@app/db/schemas"; import { BadRequestError, DatabaseError, NotFoundError } from "@app/lib/errors"; import { ormify, selectAllTableCols, sqlNestRelationships } from "@app/lib/knex"; +import { logger } from "@app/lib/logger"; +import { QueueName, TQueueServiceFactory } from "@app/queue"; export type TSecretDALFactory = ReturnType; @@ -339,6 +341,94 @@ export const secretDALFactory = (db: TDbClient) => { } }; + const pruneSecretReminders = async (queueService: TQueueServiceFactory) => { + const REMINDER_PRUNE_BATCH_SIZE = 5_000; + const MAX_RETRY_ON_FAILURE = 3; + let numberOfRetryOnFailure = 0; + let deletedReminderCount = 0; + + logger.info(`${QueueName.DailyResourceCleanUp}: secret reminders started`); + + try { + const repeatableJobs = await queueService.getRepeatableJobs(QueueName.SecretReminder); + const reminderJobs = repeatableJobs + .map((job) => ({ secretId: job.id?.replace("reminder-", "") as string, jobKey: job.key })) + .filter(Boolean); + + if (reminderJobs.length === 0) { + logger.info(`${QueueName.DailyResourceCleanUp}: no reminder jobs found`); + return; + } + + for (let offset = 0; offset < reminderJobs.length; offset += REMINDER_PRUNE_BATCH_SIZE) { + try { + const batchIds = reminderJobs.slice(offset, offset + REMINDER_PRUNE_BATCH_SIZE).map((r) => r.secretId); + + const payload = { + $in: { + id: batchIds + } + }; + + const opts = { + limit: REMINDER_PRUNE_BATCH_SIZE + }; + + // Find existing secrets with pagination + // eslint-disable-next-line no-await-in-loop + const [secrets, secretsV2] = await Promise.all([ + ormify(db, TableName.Secret).find(payload, opts), + ormify(db, TableName.SecretV2).find(payload, opts) + ]); + + const foundSecretIds = new Set([ + ...secrets.map((secret) => secret.id), + ...secretsV2.map((secret) => secret.id) + ]); + + // Find IDs that don't exist in either table + const secretIdsNotFound = batchIds.filter((secretId) => !foundSecretIds.has(secretId)); + + // Delete reminders for non-existent secrets + for (const secretId of secretIdsNotFound) { + const jobKey = reminderJobs.find((r) => r.secretId === secretId)?.jobKey; + + if (jobKey) { + // eslint-disable-next-line no-await-in-loop + await queueService.stopRepeatableJobByKey(QueueName.SecretReminder, jobKey); + deletedReminderCount += 1; + } + } + + numberOfRetryOnFailure = 0; + } catch (error) { + numberOfRetryOnFailure += 1; + logger.error(error, `Failed to process batch at offset ${offset}`); + + if (numberOfRetryOnFailure >= MAX_RETRY_ON_FAILURE) { + break; + } + + // Retry the current batch + offset -= REMINDER_PRUNE_BATCH_SIZE; + + // eslint-disable-next-line no-promise-executor-return, @typescript-eslint/no-loop-func, no-await-in-loop + await new Promise((resolve) => setTimeout(resolve, 500 * numberOfRetryOnFailure)); + } + + // Small delay between batches + // eslint-disable-next-line no-promise-executor-return, @typescript-eslint/no-loop-func, no-await-in-loop + await new Promise((resolve) => setTimeout(resolve, 10)); + } + } catch (error) { + logger.error(error, "Failed to complete secret reminder pruning"); + } finally { + logger.info( + `${QueueName.DailyResourceCleanUp}: secret reminders completed. Deleted ${deletedReminderCount} reminders` + ); + } + }; + return { ...secretOrm, update, @@ -352,6 +442,7 @@ export const secretDALFactory = (db: TDbClient) => { findByBlindIndexes, upsertSecretReferences, findReferencedSecretReferences, - findAllProjectSecretValues + findAllProjectSecretValues, + pruneSecretReminders }; };