From 4502d394a3a2b5eec593122ba42bafdd844ff448 Mon Sep 17 00:00:00 2001 From: = Date: Wed, 29 May 2024 16:29:29 +0530 Subject: [PATCH] feat: added back dedupe queue for both replication and syncing ops --- .../secret-approval-request-service.ts | 7 +- .../secret-replication-service.ts | 53 ++++-- backend/src/services/secret/secret-queue.ts | 173 ++++++++++-------- backend/src/services/secret/secret-service.ts | 15 +- backend/src/services/secret/secret-types.ts | 2 +- 5 files changed, 143 insertions(+), 107 deletions(-) diff --git a/backend/src/ee/services/secret-approval-request/secret-approval-request-service.ts b/backend/src/ee/services/secret-approval-request/secret-approval-request-service.ts index 67a477141..175435e0a 100644 --- a/backend/src/ee/services/secret-approval-request/secret-approval-request-service.ts +++ b/backend/src/ee/services/secret-approval-request/secret-approval-request-service.ts @@ -15,13 +15,13 @@ import { ActorType } from "@app/services/auth/auth-type"; import { TProjectDALFactory } from "@app/services/project/project-dal"; import { TProjectBotServiceFactory } from "@app/services/project-bot/project-bot-service"; import { TSecretDALFactory } from "@app/services/secret/secret-dal"; -import { getAllNestedSecretReferences } from "@app/services/secret/secret-fns"; import { fnSecretBlindIndexCheck, fnSecretBlindIndexCheckV2, fnSecretBulkDelete, fnSecretBulkInsert, - fnSecretBulkUpdate + fnSecretBulkUpdate, + getAllNestedSecretReferences } from "@app/services/secret/secret-fns"; import { TSecretQueueFactory } from "@app/services/secret/secret-queue"; import { SecretOperations } from "@app/services/secret/secret-types"; @@ -51,6 +51,7 @@ import { type TSecretApprovalRequestServiceFactoryDep = { permissionService: Pick; + projectBotService: Pick; secretApprovalRequestDAL: TSecretApprovalRequestDALFactory; secretApprovalRequestSecretDAL: TSecretApprovalRequestSecretDALFactory; secretApprovalRequestReviewerDAL: TSecretApprovalRequestReviewerDALFactory; @@ -356,7 +357,7 @@ export const secretApprovalRequestServiceFactory = ({ } const secretDeletionCommits = secretApprovalSecrets.filter(({ op }) => op === SecretOperations.Delete); - + const botKey = await projectBotService.getBotKey(projectId).catch(() => null); const mergeStatus = await secretApprovalRequestDAL.transaction(async (tx) => { const newSecrets = secretCreationCommits.length ? await fnSecretBulkInsert({ diff --git a/backend/src/services/secret-replication/secret-replication-service.ts b/backend/src/services/secret-replication/secret-replication-service.ts index ab1cf1c40..a289e4600 100644 --- a/backend/src/services/secret-replication/secret-replication-service.ts +++ b/backend/src/services/secret-replication/secret-replication-service.ts @@ -7,7 +7,7 @@ import { BadRequestError } from "@app/lib/errors"; import { groupBy } from "@app/lib/fn"; import { logger } from "@app/lib/logger"; import { alphaNumericNanoId } from "@app/lib/nanoid"; -import { QueueJobs, QueueName, TQueueServiceFactory } from "@app/queue"; +import { QueueName, TQueueServiceFactory } from "@app/queue"; import { ActorType } from "../auth/auth-type"; import { TProjectMembershipDALFactory } from "../project-membership/project-membership-dal"; @@ -25,7 +25,10 @@ import { TSecretReplicationDALFactory } from "./secret-replication-dal"; type TSecretReplicationServiceFactoryDep = { secretReplicationDAL: TSecretReplicationDALFactory; - secretDAL: Pick; + secretDAL: Pick< + TSecretDALFactory, + "find" | "findByBlindIndexes" | "insertMany" | "bulkUpdate" | "delete" | "upsertSecretReferences" + >; secretVersionDAL: Pick; secretImportDAL: Pick; folderDAL: Pick; @@ -48,6 +51,7 @@ type TSecretReplicationServiceFactoryDep = { export type TSecretReplicationServiceFactory = ReturnType; const SECRET_IMPORT_SUCCESS_LOCK = 10; const keystoreReplicationSuccessKey = (jobId: string, secretImportId: string) => `${jobId}-${secretImportId}`; +const getReplicationKeyLockPrefix = (keyName: string) => `REPLICATION_SECRET_${keyName}`; export const secretReplicationServiceFactory = ({ secretReplicationDAL, @@ -68,7 +72,20 @@ export const secretReplicationServiceFactory = ({ }: TSecretReplicationServiceFactoryDep) => { queueService.start(QueueName.SecretReplication, async (job) => { logger.info(job.data, "Replication started"); - const { secrets, folderId, secretPath, environmentId, projectId, actorId, actor, pickOnlyImportIds } = job.data; + const { + secrets, + folderId, + secretPath, + environmentId, + projectId, + actorId, + actor, + pickOnlyImportIds, + _deDupeReplicationQueue: deDupeReplicationQueue, + _deDupeQueue: deDupeQueue + } = job.data; + + // filter for initial filling let secretImports = await secretImportDAL.find({ importPath: secretPath, importEnv: environmentId, @@ -91,7 +108,7 @@ export const secretReplicationServiceFactory = ({ if (!sanitizedSecrets.length) return; const lock = await keyStore.acquireLock( - replicatedSecrets.map(({ id }) => id), + replicatedSecrets.map(({ id }) => getReplicationKeyLockPrefix(id)), 5000 ); @@ -301,28 +318,26 @@ export const secretReplicationServiceFactory = ({ } }); - await queueService.queue(QueueName.SecretReplication, QueueJobs.SecretReplication, { - folderId: importedFolder.id, - projectId, - secrets: nestedImportSecrets, - secretPath: importedFolder.path, - environmentId: importedFolder.envId, - actorId, - actor - }); const folderLock = await keyStore .acquireLock([`secret-replication-${importFolderId}`], 5000) .catch(() => null); if (folderLock) { await snapshotService.performSnapshot(importFolderId); await folderLock.release(); - await secretQueueService.syncSecrets({ - excludeReplication: true, - projectId, - secretPath: importedFolder.path, - environmentSlug: importedFolder.environmentSlug - }); } + + await secretQueueService.syncSecrets({ + projectId, + secretPath: importedFolder.path, + _deDupeReplicationQueue: deDupeReplicationQueue, + _deDupeQueue: deDupeQueue, + environmentSlug: importedFolder.environmentSlug, + actorId, + actor, + secrets: nestedImportSecrets, + folderId: importedFolder.id, + environmentId: importedFolder.envId + }); } await keyStore.setItemWithExpiry( diff --git a/backend/src/services/secret/secret-queue.ts b/backend/src/services/secret/secret-queue.ts index 3055c48c3..32be44bc1 100644 --- a/backend/src/services/secret/secret-queue.ts +++ b/backend/src/services/secret/secret-queue.ts @@ -64,8 +64,10 @@ export type TGetSecrets = { }; const MAX_SYNC_SECRET_DEPTH = 5; -const uniqueIntegrationKey = (environment: string, secretPath: string) => `integration-${environment}-${secretPath}`; +const uniqueSecretQueueKey = (environment: string, secretPath: string) => + `secret-queue-dedupe-${environment}-${secretPath}`; +type TIntegrationSecret = Record; export const secretQueueFactory = ({ queueService, integrationDAL, @@ -86,75 +88,6 @@ export const secretQueueFactory = ({ secretTagDAL, secretVersionTagDAL }: TSecretQueueFactoryDep) => { - const createManySecretsRawFn = createManySecretsRawFnFactory({ - projectDAL, - projectBotDAL, - secretDAL, - secretVersionDAL, - secretBlindIndexDAL, - secretTagDAL, - secretVersionTagDAL, - folderDAL - }); - - const updateManySecretsRawFn = updateManySecretsRawFnFactory({ - projectDAL, - projectBotDAL, - secretDAL, - secretVersionDAL, - secretBlindIndexDAL, - secretTagDAL, - secretVersionTagDAL, - folderDAL - }); - - const syncIntegrations = async (dto: TGetSecrets & { deDupeQueue?: Record }) => { - await queueService.queue(QueueName.IntegrationSync, QueueJobs.IntegrationSync, dto, { - attempts: 3, - delay: 1000, - backoff: { - type: "exponential", - delay: 3000 - }, - removeOnComplete: true, - removeOnFail: true - }); - }; - - const syncSecrets = async ({_deDupeQueue:deDupeQueue = {},_depth = 0, ...dto}: TSyncSecretsDTO) => { - logger.info( - `syncSecrets: syncing project secrets where [projectId=${dto.projectId}] [environment=${dto.environmentSlug}] [path=${dto.secretPath}]` - ); - const deDuplicationKey = uniqueIntegrationKey(dto.environmentSlug, dto.secretPath); - if (deDupeQueue?.[deDuplicationKey]) { - return; - } - // eslint-disable-next-line - deDupeQueue[deDuplicationKey] = true; - await queueService.queue(QueueName.SecretSync, QueueJobs.SecretSync, dto as TSyncSecretsDTO, { - removeOnFail: true, - removeOnComplete: true, - delay: 1000, - attempts: 5, - backoff: { - type: "exponential", - delay: 3000 - } - }); - }; - - const replicateSecrets = async (dto: TSyncSecretsDTO) => { - await queueService.queue(QueueName.SecretReplication, QueueJobs.SecretReplication, dto, { - attempts: 3, - backoff: { - type: "exponential", - delay: 2000 - }, - removeOnComplete: true, - removeOnFail: true - }); - }; - const removeSecretReminder = async (dto: TRemoveSecretReminderDTO) => { const appCfg = getConfig(); await queueService.stopRepeatableJob( @@ -249,8 +182,27 @@ export const secretQueueFactory = ({ } } }; + const createManySecretsRawFn = createManySecretsRawFnFactory({ + projectDAL, + projectBotDAL, + secretDAL, + secretVersionDAL, + secretBlindIndexDAL, + secretTagDAL, + secretVersionTagDAL, + folderDAL + }); - type Content = Record; + const updateManySecretsRawFn = updateManySecretsRawFnFactory({ + projectDAL, + projectBotDAL, + secretDAL, + secretVersionDAL, + secretBlindIndexDAL, + secretTagDAL, + secretVersionTagDAL, + folderDAL + }); /** * Return the secrets in a given [folderId] including secrets from @@ -263,7 +215,7 @@ export const secretQueueFactory = ({ key: string; depth: number; }) => { - let content: Content = {}; + let content: TIntegrationSecret = {}; if (dto.depth > MAX_SYNC_SECRET_DEPTH) { logger.info( `getIntegrationSecrets: secret depth exceeded for [projectId=${dto.projectId}] [folderId=${dto.folderId}] [depth=${dto.depth}]` @@ -345,8 +297,69 @@ export const secretQueueFactory = ({ return content; }; + const syncIntegrations = async (dto: TGetSecrets & { deDupeQueue?: Record }) => { + await queueService.queue(QueueName.IntegrationSync, QueueJobs.IntegrationSync, dto, { + attempts: 3, + delay: 1000, + backoff: { + type: "exponential", + delay: 3000 + }, + removeOnComplete: true, + removeOnFail: true + }); + }; + + const replicateSecrets = async (dto: Omit) => { + await queueService.queue(QueueName.SecretReplication, QueueJobs.SecretReplication, dto, { + attempts: 3, + backoff: { + type: "exponential", + delay: 2000 + }, + removeOnComplete: true, + removeOnFail: true + }); + }; + + const syncSecrets = async ({ + // seperate de-dupe queue for integration sync and replication sync + _deDupeQueue: deDupeQueue = {}, + _deDupeReplicationQueue: deDupeReplicationQueue = {}, + ...dto + }: TSyncSecretsDTO) => { + logger.info( + `syncSecrets: syncing project secrets where [projectId=${dto.projectId}] [environment=${dto.environmentSlug}] [path=${dto.secretPath}]` + ); + const deDuplicationKey = uniqueSecretQueueKey(dto.environmentSlug, dto.secretPath); + if (!dto.excludeReplication ? deDupeReplicationQueue?.[deDuplicationKey] : deDupeQueue?.[deDuplicationKey]) { + return; + } + // eslint-disable-next-line + deDupeQueue[deDuplicationKey] = true; + // eslint-disable-next-line + deDupeReplicationQueue[deDuplicationKey] = true; + await queueService.queue( + QueueName.SecretSync, + QueueJobs.SecretSync, + { ...dto, _deDupeQueue: deDupeQueue, _deDupeReplicationQueue: deDupeReplicationQueue } as TSyncSecretsDTO, + { + removeOnFail: true, + removeOnComplete: true, + delay: 1000, + attempts: 5, + backoff: { + type: "exponential", + delay: 3000 + } + } + ); + }; + queueService.start(QueueName.SecretSync, async (job) => { const { + _deDupeQueue: deDupeQueue, + _deDupeReplicationQueue: deDupeReplicationQueue, secretPath, environmentId, projectId, @@ -373,9 +386,10 @@ export const secretQueueFactory = ({ } } ); - await syncIntegrations({ secretPath, projectId, environment }); - if (!excludeReplication) + await syncIntegrations({ secretPath, projectId, environment, deDupeQueue }); + if (!excludeReplication) { await replicateSecrets({ + _deDupeReplicationQueue: deDupeReplicationQueue, environmentId, projectId, secretPath, @@ -386,6 +400,7 @@ export const secretQueueFactory = ({ excludeReplication, environmentSlug: environment }); + } }); queueService.start(QueueName.IntegrationSync, async (job) => { @@ -409,7 +424,7 @@ export const secretQueueFactory = ({ const imports = await secretImportDAL.find(linkSourceDto); if (imports.length) { - // keep calling sync secret for all the imports made + // keep calling sync secret for all the imports made const importedFolderIds = unique(imports, (i) => i.folderId).map(({ folderId }) => folderId); const importedFolders = await folderDAL.findSecretPathByFolderIds(projectId, importedFolderIds); const foldersGroupedById = groupBy(importedFolders.filter(Boolean), (i) => i?.id as string); @@ -423,15 +438,14 @@ export const secretQueueFactory = ({ .filter( ({ folderId }) => !deDupeQueue[ - uniqueIntegrationKey( + uniqueSecretQueueKey( foldersGroupedById[folderId][0]?.environmentSlug as string, - foldersGroupedById[folderId][0]?.path + foldersGroupedById[folderId][0]?.path as string ) ] ) .map(({ folderId }) => syncSecrets({ - _depth: depth + 1, projectId, secretPath: foldersGroupedById[folderId][0]?.path as string, environmentSlug: foldersGroupedById[folderId][0]?.environmentSlug as string, @@ -461,7 +475,7 @@ export const secretQueueFactory = ({ .filter( ({ folderId }) => !deDupeQueue[ - uniqueIntegrationKey( + uniqueSecretQueueKey( referencedFoldersGroupedById[folderId][0]?.environmentSlug as string, referencedFoldersGroupedById[folderId][0]?.path as string ) @@ -469,7 +483,6 @@ export const secretQueueFactory = ({ ) .map(({ folderId }) => syncSecrets({ - _depth: depth + 1, projectId, secretPath: referencedFoldersGroupedById[folderId][0]?.path as string, environmentSlug: referencedFoldersGroupedById[folderId][0]?.environmentSlug as string, diff --git a/backend/src/services/secret/secret-service.ts b/backend/src/services/secret/secret-service.ts index 35bff3772..d43b48268 100644 --- a/backend/src/services/secret/secret-service.ts +++ b/backend/src/services/secret/secret-service.ts @@ -1251,7 +1251,9 @@ export const secretServiceFactory = ({ }) }); - return secrets.map((secret) => decryptSecretRaw({ ...secret, workspace: projectId, environment }, botKey)); + return secrets.map((secret) => + decryptSecretRaw({ ...secret, workspace: projectId, environment, secretPath }, botKey) + ); }; const updateManySecretsRaw = async ({ @@ -1300,7 +1302,9 @@ export const secretServiceFactory = ({ }) }); - return secrets.map((secret) => decryptSecretRaw({ ...secret, workspace: projectId, environment }, botKey)); + return secrets.map((secret) => + decryptSecretRaw({ ...secret, workspace: projectId, environment, secretPath }, botKey) + ); }; const deleteManySecretsRaw = async ({ @@ -1331,7 +1335,9 @@ export const secretServiceFactory = ({ secrets: inputSecrets.map(({ secretKey }) => ({ secretName: secretKey, type: SecretType.Shared })) }); - return secrets.map((secret) => decryptSecretRaw({ ...secret, workspace: projectId, environment }, botKey)); + return secrets.map((secret) => + decryptSecretRaw({ ...secret, workspace: projectId, environment, secretPath }, botKey) + ); }; const getSecretVersions = async ({ @@ -1637,6 +1643,7 @@ export const secretServiceFactory = ({ createManySecretsRaw, updateManySecretsRaw, deleteManySecretsRaw, - getSecretVersions + getSecretVersions, + backfillSecretReferences }; }; diff --git a/backend/src/services/secret/secret-types.ts b/backend/src/services/secret/secret-types.ts index fb62ec052..8c6d5db06 100644 --- a/backend/src/services/secret/secret-types.ts +++ b/backend/src/services/secret/secret-types.ts @@ -378,8 +378,8 @@ export enum SecretOperations { } export type TSyncSecretsDTO = { - _depth?: number; _deDupeQueue?: Record; + _deDupeReplicationQueue?: Record; secretPath: string; projectId: string; environmentSlug: string;