From 3541ddf8ac8ccbd43cf0d9e1ab10db3255db8944 Mon Sep 17 00:00:00 2001 From: Akhil Mohan Date: Thu, 2 May 2024 18:04:36 +0530 Subject: [PATCH] feat: fixed folder dal wrong folder details in findBySecretPath issue and replication dal --- .../secret-folder/secret-folder-dal.ts | 17 ++++- .../secret-replication-dal.ts | 23 ++---- .../secret-replication-service.ts | 72 +++++++++++++------ 3 files changed, 69 insertions(+), 43 deletions(-) diff --git a/backend/src/services/secret-folder/secret-folder-dal.ts b/backend/src/services/secret-folder/secret-folder-dal.ts index b3147d1fa..0e896d0c6 100644 --- a/backend/src/services/secret-folder/secret-folder-dal.ts +++ b/backend/src/services/secret-folder/secret-folder-dal.ts @@ -169,6 +169,7 @@ const sqlFindSecretPathByFolderId = (db: Knex, projectId: string, folderIds: str // this is for root condition // if the given folder id is root folder id then intial path is set as / instead of /root // if not root folder the path here will be / + depth: 1, path: db.raw(`CONCAT('/', (CASE WHEN "parentId" is NULL THEN '' ELSE ${TableName.SecretFolder}.name END))`), child: db.raw("NULL::uuid"), environmentSlug: `${TableName.Environment}.slug` @@ -185,6 +186,7 @@ const sqlFindSecretPathByFolderId = (db: Knex, projectId: string, folderIds: str .select({ // then we join join this folder name behind previous as we are going from child to parent // the root folder check is used to avoid last / and also root name in folders + depth: db.raw("parent.depth + 1"), path: db.raw( `CONCAT( CASE WHEN ${TableName.SecretFolder}."parentId" is NULL THEN '' @@ -199,7 +201,7 @@ const sqlFindSecretPathByFolderId = (db: Knex, projectId: string, folderIds: str ); }) .select("*") - .from("parent"); + .from("parent"); export type TSecretFolderDALFactory = ReturnType; // never change this. If u do write a migration for it @@ -260,12 +262,23 @@ export const secretFolderDALFactory = (db: TDbClient) => { try { const folders = await sqlFindSecretPathByFolderId(tx || db, projectId, folderIds); + // travelling all the way from leaf node to root contains real path const rootFolders = groupBy( folders.filter(({ parentId }) => parentId === null), (i) => i.child || i.id // root condition then child and parent will null ); + const actualFolders = groupBy( + folders.filter(({ depth }) => depth === 1), + (i) => i.id // root condition then child and parent will null + ); - return folderIds.map((folderId) => rootFolders[folderId]?.[0]); + return folderIds.map((folderId) => { + if (!rootFolders[folderId]?.[0]) return; + + const actualId = rootFolders[folderId][0].child || rootFolders[folderId][0].id; + const folder = actualFolders[actualId][0]; + return { ...folder, path: rootFolders[folderId]?.[0].path }; + }); } catch (error) { throw new DatabaseError({ error, name: "Find by secret path" }); } diff --git a/backend/src/services/secret-replication/secret-replication-dal.ts b/backend/src/services/secret-replication/secret-replication-dal.ts index d977f3fa2..2a8b035ea 100644 --- a/backend/src/services/secret-replication/secret-replication-dal.ts +++ b/backend/src/services/secret-replication/secret-replication-dal.ts @@ -25,27 +25,14 @@ export const secretReplicationDALFactory = (db: TDbClient) => { .leftJoin( (tx || db)(TableName.SecretVersion) .where("isReplicated", true) - .groupBy(["secretId", "version"]) + .groupBy("secretId") .max("version") - .select("version", "secretId") + .select("secretId") .as("latestVersion"), - (bd) => { - bd.on(`${TableName.SecretVersion}.secretId`, "latestVersion.secretId").andOn( - `${TableName.SecretVersion}.version`, - "latestVersion.max" - ); - } + `${TableName.SecretVersion}.secretId`, + "latestVersion.secretId" ) - // .leftJoin( - // (tx || db)(TableName.SecretVersion).select("isReplicated", "version", "secretId").as("previousVersion"), - // (bd) => { - // bd.on(`${TableName.SecretVersion}.secretId`, "previousVersion.secretId").andOn( - // "previousVersion.version", - // (tx || db).raw(`${TableName.SecretVersion}.version - 1`) - // ); - // } - // ) - .select(db.ref("version").withSchema("latestVersion").as("latestReplicatedVersion")) + .select(db.ref("max").withSchema("latestVersion").as("latestReplicatedVersion")) .select(selectAllTableCols(TableName.SecretVersion)); return sqlRawDocs; diff --git a/backend/src/services/secret-replication/secret-replication-service.ts b/backend/src/services/secret-replication/secret-replication-service.ts index 3b8c12f58..acb290254 100644 --- a/backend/src/services/secret-replication/secret-replication-service.ts +++ b/backend/src/services/secret-replication/secret-replication-service.ts @@ -1,7 +1,9 @@ import { TSecretApprovalPolicyServiceFactory } from "@app/ee/services/secret-approval-policy/secret-approval-policy-service"; import { TSecretApprovalRequestDALFactory } from "@app/ee/services/secret-approval-request/secret-approval-request-dal"; import { TSecretApprovalRequestSecretDALFactory } from "@app/ee/services/secret-approval-request/secret-approval-request-secret-dal"; -import { TKeyStoreFactory } from "@app/keystore/keystore"; +import { TSecretSnapshotServiceFactory } from "@app/ee/services/secret-snapshot/secret-snapshot-service"; +import { KeyStorePrefixes, TKeyStoreFactory } from "@app/keystore/keystore"; +import { BadRequestError } from "@app/lib/errors"; import { groupBy } from "@app/lib/fn"; import { logger } from "@app/lib/logger"; import { alphaNumericNanoId } from "@app/lib/nanoid"; @@ -9,6 +11,7 @@ import { QueueJobs, QueueName, TQueueServiceFactory } from "@app/queue"; import { TSecretDALFactory } from "../secret/secret-dal"; import { fnSecretBulkInsert, fnSecretBulkUpdate } from "../secret/secret-fns"; +import { TSecretQueueFactory } from "../secret/secret-queue"; import { SecretOperations, TSyncSecretsDTO } from "../secret/secret-types"; import { TSecretVersionDALFactory } from "../secret/secret-version-dal"; import { TSecretVersionTagDALFactory } from "../secret/secret-version-tag-dal"; @@ -25,9 +28,11 @@ type TSecretReplicationServiceFactoryDep = { secretImportDAL: Pick; folderDAL: Pick; secretVersionTagDAL: Pick; + secretQueueService: Pick; + snapshotService: Pick; queueService: Pick; secretApprovalPolicyService: Pick; - keyStore: Pick; + keyStore: Pick; secretBlindIndexDAL: Pick; secretTagDAL: Pick; secretApprovalRequestDAL: Pick; @@ -38,14 +43,8 @@ type TSecretReplicationServiceFactoryDep = { }; export type TSecretReplicationServiceFactory = ReturnType; - -// function getRandomError(): number { -// const minCeiled: number = Math.ceil(0); -// const maxFloored: number = Math.floor(20); -// const val = Math.floor(Math.random() * (maxFloored - minCeiled) + minCeiled); // The maximum is exclusive and the minimum is inclusive -// if (val >= 10) throw new Error("Random error point"); -// return val; -// } +const SECRET_IMPORT_SUCCESS_LOCK = 10; +const keystoreReplicationSuccessKey = (jobId: string, secretImportId: string) => `${jobId}-${secretImportId}`; export const secretReplicationServiceFactory = ({ secretReplicationDAL, @@ -59,7 +58,9 @@ export const secretReplicationServiceFactory = ({ folderDAL, secretApprovalPolicyService, secretApprovalRequestSecretDAL, - secretApprovalRequestDAL + secretApprovalRequestDAL, + secretQueueService, + snapshotService }: TSecretReplicationServiceFactoryDep) => { queueService.start(QueueName.SecretReplication, async (job) => { logger.info(job.data, "Replication started"); @@ -69,11 +70,9 @@ export const secretReplicationServiceFactory = ({ importEnv: environmentId, isReplication: true }); - console.log(">>>> Secret Imports replics ", secretImports.length, secretPath, environmentId); if (!secretImports.length || !secrets.length) return; // unfiltered secrets to be replicated - console.log(secrets.length); const toBeReplicatedSecrets = await secretReplicationDAL.findSecrets({ folderId, secrets }); const replicatedSecrets = toBeReplicatedSecrets.filter( ({ version, latestReplicatedVersion, secretBlindIndex }) => @@ -81,7 +80,6 @@ export const secretReplicationServiceFactory = ({ ); const replicatedSecretsGroupBySecretId = groupBy(replicatedSecrets, (i) => i.secretId); - console.log("replicated ", replicatedSecretsGroupBySecretId); const lock = await keyStore.acquireLock( replicatedSecrets.map(({ id }) => id), 5000 @@ -90,13 +88,26 @@ export const secretReplicationServiceFactory = ({ try { /* eslint-disable no-await-in-loop */ for (const secretImport of secretImports) { + const hasJobCompleted = await keyStore.getItem( + keystoreReplicationSuccessKey(job.id as string, secretImport.id), + KeyStorePrefixes.SecretReplication + ); + if (hasJobCompleted) { + logger.info( + { jobId: job.id, importId: secretImport.id }, + "Skipping this job as this has been successfully replicated." + ); + // eslint-disable-next-line + continue; + } + const [importedFolder] = await folderDAL.findSecretPathByFolderIds(projectId, [secretImport.folderId]); + if (!importedFolder) throw new BadRequestError({ message: "Imported folder not found" }); const importFolderId = importedFolder.id; const localSecrets = await secretDAL.find({ $in: { secretBlindIndex: replicatedSecrets.map(({ secretBlindIndex }) => secretBlindIndex) }, - folderId: importFolderId, - isReplicated: true + folderId: importFolderId }); const localSecretsGroupedByBlindIndex = groupBy(localSecrets, (i) => i.secretBlindIndex as string); @@ -113,11 +124,6 @@ export const secretReplicationServiceFactory = ({ localSecretsGroupedByBlindIndex[replicatedSecretsGroupBySecretId[id][0].secretBlindIndex as string]?.[0] ); - console.log(replicatedSecretsGroupBySecretId); - console.log(locallyCreatedSecrets); - console.log("update", locallyUpdatedSecrets); - console.log("local board", localSecrets); - const locallyDeletedSecrets = secrets.filter( ({ operation, id }) => operation === SecretOperations.Delete && @@ -277,7 +283,7 @@ export const secretReplicationServiceFactory = ({ ); } }); - console.log("Environment ID -> slug", importedFolder.envId, importedFolder.environmentSlug); + await queueService.queue(QueueName.SecretReplication, QueueJobs.SecretReplication, { folderId: importedFolder.id, projectId, @@ -286,7 +292,27 @@ export const secretReplicationServiceFactory = ({ environmentId: importedFolder.envId, membershipId }); + 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 keyStore.setItemWithExpiry( + keystoreReplicationSuccessKey(job.id as string, secretImport.id), + SECRET_IMPORT_SUCCESS_LOCK, + 1, + KeyStorePrefixes.SecretReplication + ); } await secretVersionDAL.update({ $in: { id: replicatedSecrets.map(({ id }) => id) } }, { isReplicated: true }); /* eslint-enable no-await-in-loop */ @@ -295,7 +321,7 @@ export const secretReplicationServiceFactory = ({ } }); - queueService.listen(QueueName.SecretReplication, "failed", async (job, err) => { + queueService.listen(QueueName.SecretReplication, "failed", (job, err) => { logger.error(err, "Failed to replicate secret", job?.data); }); };