feat: added back dedupe queue for both replication and syncing ops

This commit is contained in:
=
2024-05-31 23:10:02 +05:30
parent 531d3751a8
commit 4502d394a3
5 changed files with 143 additions and 107 deletions
@@ -15,13 +15,13 @@ import { ActorType } from "@app/services/auth/auth-type";
import { TProjectDALFactory } from "@app/services/project/project-dal"; import { TProjectDALFactory } from "@app/services/project/project-dal";
import { TProjectBotServiceFactory } from "@app/services/project-bot/project-bot-service"; import { TProjectBotServiceFactory } from "@app/services/project-bot/project-bot-service";
import { TSecretDALFactory } from "@app/services/secret/secret-dal"; import { TSecretDALFactory } from "@app/services/secret/secret-dal";
import { getAllNestedSecretReferences } from "@app/services/secret/secret-fns";
import { import {
fnSecretBlindIndexCheck, fnSecretBlindIndexCheck,
fnSecretBlindIndexCheckV2, fnSecretBlindIndexCheckV2,
fnSecretBulkDelete, fnSecretBulkDelete,
fnSecretBulkInsert, fnSecretBulkInsert,
fnSecretBulkUpdate fnSecretBulkUpdate,
getAllNestedSecretReferences
} from "@app/services/secret/secret-fns"; } from "@app/services/secret/secret-fns";
import { TSecretQueueFactory } from "@app/services/secret/secret-queue"; import { TSecretQueueFactory } from "@app/services/secret/secret-queue";
import { SecretOperations } from "@app/services/secret/secret-types"; import { SecretOperations } from "@app/services/secret/secret-types";
@@ -51,6 +51,7 @@ import {
type TSecretApprovalRequestServiceFactoryDep = { type TSecretApprovalRequestServiceFactoryDep = {
permissionService: Pick<TPermissionServiceFactory, "getProjectPermission">; permissionService: Pick<TPermissionServiceFactory, "getProjectPermission">;
projectBotService: Pick<TProjectBotServiceFactory, "getBotKey">;
secretApprovalRequestDAL: TSecretApprovalRequestDALFactory; secretApprovalRequestDAL: TSecretApprovalRequestDALFactory;
secretApprovalRequestSecretDAL: TSecretApprovalRequestSecretDALFactory; secretApprovalRequestSecretDAL: TSecretApprovalRequestSecretDALFactory;
secretApprovalRequestReviewerDAL: TSecretApprovalRequestReviewerDALFactory; secretApprovalRequestReviewerDAL: TSecretApprovalRequestReviewerDALFactory;
@@ -356,7 +357,7 @@ export const secretApprovalRequestServiceFactory = ({
} }
const secretDeletionCommits = secretApprovalSecrets.filter(({ op }) => op === SecretOperations.Delete); const secretDeletionCommits = secretApprovalSecrets.filter(({ op }) => op === SecretOperations.Delete);
const botKey = await projectBotService.getBotKey(projectId).catch(() => null);
const mergeStatus = await secretApprovalRequestDAL.transaction(async (tx) => { const mergeStatus = await secretApprovalRequestDAL.transaction(async (tx) => {
const newSecrets = secretCreationCommits.length const newSecrets = secretCreationCommits.length
? await fnSecretBulkInsert({ ? await fnSecretBulkInsert({
@@ -7,7 +7,7 @@ import { BadRequestError } from "@app/lib/errors";
import { groupBy } from "@app/lib/fn"; import { groupBy } from "@app/lib/fn";
import { logger } from "@app/lib/logger"; import { logger } from "@app/lib/logger";
import { alphaNumericNanoId } from "@app/lib/nanoid"; 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 { ActorType } from "../auth/auth-type";
import { TProjectMembershipDALFactory } from "../project-membership/project-membership-dal"; import { TProjectMembershipDALFactory } from "../project-membership/project-membership-dal";
@@ -25,7 +25,10 @@ import { TSecretReplicationDALFactory } from "./secret-replication-dal";
type TSecretReplicationServiceFactoryDep = { type TSecretReplicationServiceFactoryDep = {
secretReplicationDAL: TSecretReplicationDALFactory; secretReplicationDAL: TSecretReplicationDALFactory;
secretDAL: Pick<TSecretDALFactory, "find" | "findByBlindIndexes" | "insertMany" | "bulkUpdate" | "delete">; secretDAL: Pick<
TSecretDALFactory,
"find" | "findByBlindIndexes" | "insertMany" | "bulkUpdate" | "delete" | "upsertSecretReferences"
>;
secretVersionDAL: Pick<TSecretVersionDALFactory, "find" | "insertMany" | "update" | "findLatestVersionMany">; secretVersionDAL: Pick<TSecretVersionDALFactory, "find" | "insertMany" | "update" | "findLatestVersionMany">;
secretImportDAL: Pick<TSecretImportDALFactory, "find">; secretImportDAL: Pick<TSecretImportDALFactory, "find">;
folderDAL: Pick<TSecretFolderDALFactory, "findSecretPathByFolderIds" | "findBySecretPath">; folderDAL: Pick<TSecretFolderDALFactory, "findSecretPathByFolderIds" | "findBySecretPath">;
@@ -48,6 +51,7 @@ type TSecretReplicationServiceFactoryDep = {
export type TSecretReplicationServiceFactory = ReturnType<typeof secretReplicationServiceFactory>; export type TSecretReplicationServiceFactory = ReturnType<typeof secretReplicationServiceFactory>;
const SECRET_IMPORT_SUCCESS_LOCK = 10; const SECRET_IMPORT_SUCCESS_LOCK = 10;
const keystoreReplicationSuccessKey = (jobId: string, secretImportId: string) => `${jobId}-${secretImportId}`; const keystoreReplicationSuccessKey = (jobId: string, secretImportId: string) => `${jobId}-${secretImportId}`;
const getReplicationKeyLockPrefix = (keyName: string) => `REPLICATION_SECRET_${keyName}`;
export const secretReplicationServiceFactory = ({ export const secretReplicationServiceFactory = ({
secretReplicationDAL, secretReplicationDAL,
@@ -68,7 +72,20 @@ export const secretReplicationServiceFactory = ({
}: TSecretReplicationServiceFactoryDep) => { }: TSecretReplicationServiceFactoryDep) => {
queueService.start(QueueName.SecretReplication, async (job) => { queueService.start(QueueName.SecretReplication, async (job) => {
logger.info(job.data, "Replication started"); 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({ let secretImports = await secretImportDAL.find({
importPath: secretPath, importPath: secretPath,
importEnv: environmentId, importEnv: environmentId,
@@ -91,7 +108,7 @@ export const secretReplicationServiceFactory = ({
if (!sanitizedSecrets.length) return; if (!sanitizedSecrets.length) return;
const lock = await keyStore.acquireLock( const lock = await keyStore.acquireLock(
replicatedSecrets.map(({ id }) => id), replicatedSecrets.map(({ id }) => getReplicationKeyLockPrefix(id)),
5000 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 const folderLock = await keyStore
.acquireLock([`secret-replication-${importFolderId}`], 5000) .acquireLock([`secret-replication-${importFolderId}`], 5000)
.catch(() => null); .catch(() => null);
if (folderLock) { if (folderLock) {
await snapshotService.performSnapshot(importFolderId); await snapshotService.performSnapshot(importFolderId);
await folderLock.release(); 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( await keyStore.setItemWithExpiry(
+93 -80
View File
@@ -64,8 +64,10 @@ export type TGetSecrets = {
}; };
const MAX_SYNC_SECRET_DEPTH = 5; 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<string, { value: string; comment?: string; skipMultilineEncoding?: boolean }>;
export const secretQueueFactory = ({ export const secretQueueFactory = ({
queueService, queueService,
integrationDAL, integrationDAL,
@@ -86,75 +88,6 @@ export const secretQueueFactory = ({
secretTagDAL, secretTagDAL,
secretVersionTagDAL secretVersionTagDAL
}: TSecretQueueFactoryDep) => { }: 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<string, boolean> }) => {
await queueService.queue(QueueName.IntegrationSync, QueueJobs.IntegrationSync, dto, {
attempts: 3,
delay: 1000,
backoff: {
type: "exponential",
delay: 3000
},
removeOnComplete: true,
removeOnFail: true
});
};
const syncSecrets = async <T extends boolean = false>({_deDupeQueue:deDupeQueue = {},_depth = 0, ...dto}: TSyncSecretsDTO<T>) => {
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 removeSecretReminder = async (dto: TRemoveSecretReminderDTO) => {
const appCfg = getConfig(); const appCfg = getConfig();
await queueService.stopRepeatableJob( await queueService.stopRepeatableJob(
@@ -249,8 +182,27 @@ export const secretQueueFactory = ({
} }
} }
}; };
const createManySecretsRawFn = createManySecretsRawFnFactory({
projectDAL,
projectBotDAL,
secretDAL,
secretVersionDAL,
secretBlindIndexDAL,
secretTagDAL,
secretVersionTagDAL,
folderDAL
});
type Content = Record<string, { value: string; comment?: string; skipMultilineEncoding?: boolean }>; const updateManySecretsRawFn = updateManySecretsRawFnFactory({
projectDAL,
projectBotDAL,
secretDAL,
secretVersionDAL,
secretBlindIndexDAL,
secretTagDAL,
secretVersionTagDAL,
folderDAL
});
/** /**
* Return the secrets in a given [folderId] including secrets from * Return the secrets in a given [folderId] including secrets from
@@ -263,7 +215,7 @@ export const secretQueueFactory = ({
key: string; key: string;
depth: number; depth: number;
}) => { }) => {
let content: Content = {}; let content: TIntegrationSecret = {};
if (dto.depth > MAX_SYNC_SECRET_DEPTH) { if (dto.depth > MAX_SYNC_SECRET_DEPTH) {
logger.info( logger.info(
`getIntegrationSecrets: secret depth exceeded for [projectId=${dto.projectId}] [folderId=${dto.folderId}] [depth=${dto.depth}]` `getIntegrationSecrets: secret depth exceeded for [projectId=${dto.projectId}] [folderId=${dto.folderId}] [depth=${dto.depth}]`
@@ -345,8 +297,69 @@ export const secretQueueFactory = ({
return content; return content;
}; };
const syncIntegrations = async (dto: TGetSecrets & { deDupeQueue?: Record<string, boolean> }) => {
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<TSyncSecretsDTO, "deDupeQueue">) => {
await queueService.queue(QueueName.SecretReplication, QueueJobs.SecretReplication, dto, {
attempts: 3,
backoff: {
type: "exponential",
delay: 2000
},
removeOnComplete: true,
removeOnFail: true
});
};
const syncSecrets = async <T extends boolean = false>({
// seperate de-dupe queue for integration sync and replication sync
_deDupeQueue: deDupeQueue = {},
_deDupeReplicationQueue: deDupeReplicationQueue = {},
...dto
}: TSyncSecretsDTO<T>) => {
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) => { queueService.start(QueueName.SecretSync, async (job) => {
const { const {
_deDupeQueue: deDupeQueue,
_deDupeReplicationQueue: deDupeReplicationQueue,
secretPath, secretPath,
environmentId, environmentId,
projectId, projectId,
@@ -373,9 +386,10 @@ export const secretQueueFactory = ({
} }
} }
); );
await syncIntegrations({ secretPath, projectId, environment }); await syncIntegrations({ secretPath, projectId, environment, deDupeQueue });
if (!excludeReplication) if (!excludeReplication) {
await replicateSecrets({ await replicateSecrets({
_deDupeReplicationQueue: deDupeReplicationQueue,
environmentId, environmentId,
projectId, projectId,
secretPath, secretPath,
@@ -386,6 +400,7 @@ export const secretQueueFactory = ({
excludeReplication, excludeReplication,
environmentSlug: environment environmentSlug: environment
}); });
}
}); });
queueService.start(QueueName.IntegrationSync, async (job) => { queueService.start(QueueName.IntegrationSync, async (job) => {
@@ -409,7 +424,7 @@ export const secretQueueFactory = ({
const imports = await secretImportDAL.find(linkSourceDto); const imports = await secretImportDAL.find(linkSourceDto);
if (imports.length) { 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 importedFolderIds = unique(imports, (i) => i.folderId).map(({ folderId }) => folderId);
const importedFolders = await folderDAL.findSecretPathByFolderIds(projectId, importedFolderIds); const importedFolders = await folderDAL.findSecretPathByFolderIds(projectId, importedFolderIds);
const foldersGroupedById = groupBy(importedFolders.filter(Boolean), (i) => i?.id as string); const foldersGroupedById = groupBy(importedFolders.filter(Boolean), (i) => i?.id as string);
@@ -423,15 +438,14 @@ export const secretQueueFactory = ({
.filter( .filter(
({ folderId }) => ({ folderId }) =>
!deDupeQueue[ !deDupeQueue[
uniqueIntegrationKey( uniqueSecretQueueKey(
foldersGroupedById[folderId][0]?.environmentSlug as string, foldersGroupedById[folderId][0]?.environmentSlug as string,
foldersGroupedById[folderId][0]?.path foldersGroupedById[folderId][0]?.path as string
) )
] ]
) )
.map(({ folderId }) => .map(({ folderId }) =>
syncSecrets({ syncSecrets({
_depth: depth + 1,
projectId, projectId,
secretPath: foldersGroupedById[folderId][0]?.path as string, secretPath: foldersGroupedById[folderId][0]?.path as string,
environmentSlug: foldersGroupedById[folderId][0]?.environmentSlug as string, environmentSlug: foldersGroupedById[folderId][0]?.environmentSlug as string,
@@ -461,7 +475,7 @@ export const secretQueueFactory = ({
.filter( .filter(
({ folderId }) => ({ folderId }) =>
!deDupeQueue[ !deDupeQueue[
uniqueIntegrationKey( uniqueSecretQueueKey(
referencedFoldersGroupedById[folderId][0]?.environmentSlug as string, referencedFoldersGroupedById[folderId][0]?.environmentSlug as string,
referencedFoldersGroupedById[folderId][0]?.path as string referencedFoldersGroupedById[folderId][0]?.path as string
) )
@@ -469,7 +483,6 @@ export const secretQueueFactory = ({
) )
.map(({ folderId }) => .map(({ folderId }) =>
syncSecrets({ syncSecrets({
_depth: depth + 1,
projectId, projectId,
secretPath: referencedFoldersGroupedById[folderId][0]?.path as string, secretPath: referencedFoldersGroupedById[folderId][0]?.path as string,
environmentSlug: referencedFoldersGroupedById[folderId][0]?.environmentSlug as string, environmentSlug: referencedFoldersGroupedById[folderId][0]?.environmentSlug as string,
+11 -4
View File
@@ -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 ({ 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 ({ const deleteManySecretsRaw = async ({
@@ -1331,7 +1335,9 @@ export const secretServiceFactory = ({
secrets: inputSecrets.map(({ secretKey }) => ({ secretName: secretKey, type: SecretType.Shared })) 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 ({ const getSecretVersions = async ({
@@ -1637,6 +1643,7 @@ export const secretServiceFactory = ({
createManySecretsRaw, createManySecretsRaw,
updateManySecretsRaw, updateManySecretsRaw,
deleteManySecretsRaw, deleteManySecretsRaw,
getSecretVersions getSecretVersions,
backfillSecretReferences
}; };
}; };
+1 -1
View File
@@ -378,8 +378,8 @@ export enum SecretOperations {
} }
export type TSyncSecretsDTO<T extends boolean = false> = { export type TSyncSecretsDTO<T extends boolean = false> = {
_depth?: number;
_deDupeQueue?: Record<string, boolean>; _deDupeQueue?: Record<string, boolean>;
_deDupeReplicationQueue?: Record<string, boolean>;
secretPath: string; secretPath: string;
projectId: string; projectId: string;
environmentSlug: string; environmentSlug: string;