diff --git a/backend/src/keystore/keystore.ts b/backend/src/keystore/keystore.ts index 3b19012e7..bff0ad885 100644 --- a/backend/src/keystore/keystore.ts +++ b/backend/src/keystore/keystore.ts @@ -22,12 +22,15 @@ export const KeyStorePrefixes = { `sync-integration-last-run-${projectId}-${environmentSlug}-${secretPath}` as const, IdentityAccessTokenStatusUpdate: (identityAccessTokenId: string) => `identity-access-token-status:${identityAccessTokenId}`, - ServiceTokenStatusUpdate: (serviceTokenId: string) => `service-token-status:${serviceTokenId}` + ServiceTokenStatusUpdate: (serviceTokenId: string) => `service-token-status:${serviceTokenId}`, + SendFailedIntegrationSyncEmails: (projectId: string, secretPath: string, environmentSlug: string) => + `send-failed-integration-sync-emails-${projectId}-${secretPath}-${environmentSlug}` as const }; export const KeyStoreTtls = { SetSyncSecretIntegrationLastRunTimestampInSeconds: 10, - AccessTokenStatusUpdateInSeconds: 120 + AccessTokenStatusUpdateInSeconds: 120, + SendFailedIntegrationSyncEmailsInSeconds: 60 }; type TWaitTillReady = { diff --git a/backend/src/queue/queue-service.ts b/backend/src/queue/queue-service.ts index 68fcb0db2..1a5560d3a 100644 --- a/backend/src/queue/queue-service.ts +++ b/backend/src/queue/queue-service.ts @@ -7,7 +7,7 @@ import { TScanFullRepoEventPayload, TScanPushEventPayload } from "@app/ee/services/secret-scanning/secret-scanning-queue/secret-scanning-queue-types"; -import { TSyncSecretsDTO } from "@app/services/secret/secret-types"; +import { TIntegrationSyncPayload, TSyncSecretsDTO } from "@app/services/secret/secret-types"; export enum QueueName { SecretRotation = "secret-rotation", @@ -42,6 +42,7 @@ export enum QueueJobs { SecWebhook = "secret-webhook-trigger", TelemetryInstanceStats = "telemetry-self-hosted-stats", IntegrationSync = "secret-integration-pull", + SendFailedIntegrationSyncEmails = "send-failed-integration-sync-emails", SecretScan = "secret-scan", UpgradeProjectToGhost = "upgrade-project-to-ghost-job", DynamicSecretRevocation = "dynamic-secret-revocation", @@ -88,18 +89,30 @@ export type TQueueJobTypes = { name: QueueJobs.SecWebhook; payload: { projectId: string; environment: string; secretPath: string; depth?: number }; }; - [QueueName.IntegrationSync]: { - name: QueueJobs.IntegrationSync; - payload: { - isManual?: boolean; - actorId?: string; - projectId: string; - environment: string; - secretPath: string; - depth?: number; - deDupeQueue?: Record; - }; - }; + + [QueueName.AccessTokenStatusUpdate]: + | { + name: QueueJobs.IdentityAccessTokenStatusUpdate; + payload: { identityAccessTokenId: string; numberOfUses: number }; + } + | { + name: QueueJobs.ServiceTokenStatusUpdate; + payload: { serviceTokenId: string }; + }; + + [QueueName.IntegrationSync]: + | { + name: QueueJobs.IntegrationSync; + payload: TIntegrationSyncPayload; + } + | { + name: QueueJobs.SendFailedIntegrationSyncEmails; + payload: { + projectId: string; + environmentSlug: string; + secretPath: string; + }; + }; [QueueName.SecretFullRepoScan]: { name: QueueJobs.SecretScan; payload: TScanFullRepoEventPayload; @@ -153,15 +166,6 @@ export type TQueueJobTypes = { name: QueueJobs.ProjectV3Migration; payload: { projectId: string }; }; - [QueueName.AccessTokenStatusUpdate]: - | { - name: QueueJobs.IdentityAccessTokenStatusUpdate; - payload: { identityAccessTokenId: string; numberOfUses: number }; - } - | { - name: QueueJobs.ServiceTokenStatusUpdate; - payload: { serviceTokenId: string }; - }; }; export type TQueueServiceFactory = ReturnType; diff --git a/backend/src/services/secret/secret-queue.ts b/backend/src/services/secret/secret-queue.ts index 1926be6da..e858a08a4 100644 --- a/backend/src/services/secret/secret-queue.ts +++ b/backend/src/services/secret/secret-queue.ts @@ -17,7 +17,7 @@ import { TSnapshotSecretV2DALFactory } from "@app/ee/services/secret-snapshot/sn import { KeyStorePrefixes, KeyStoreTtls, TKeyStoreFactory } from "@app/keystore/keystore"; import { getConfig } from "@app/lib/config/env"; import { decryptSymmetric128BitHexKeyUTF8 } from "@app/lib/crypto"; -import { daysToMillisecond, secondsToMillis } from "@app/lib/dates"; +import { applyJitter, daysToMillisecond, secondsToMillis } from "@app/lib/dates"; import { BadRequestError } from "@app/lib/errors"; import { getTimeDifferenceInSeconds, groupBy, isSamePath, unique } from "@app/lib/fn"; import { logger } from "@app/lib/logger"; @@ -55,8 +55,11 @@ import { fnTriggerWebhook } from "../webhook/webhook-fns"; import { TSecretDALFactory } from "./secret-dal"; import { interpolateSecrets } from "./secret-fns"; import { + FailedIntegrationSyncEmailsPayloadSchema, TCreateSecretReminderDTO, + TFailedIntegrationSyncEmailsPayload, THandleReminderDTO, + TIntegrationSyncPayload, TRemoveSecretReminderDTO, TSyncSecretsDTO } from "./secret-types"; @@ -515,6 +518,38 @@ export const secretQueueFactory = ({ ); }; + const sendFailedIntegrationSyncEmails = async (payload: TFailedIntegrationSyncEmailsPayload) => { + const key = KeyStorePrefixes.SendFailedIntegrationSyncEmails( + payload.projectId, + payload.secretPath, + payload.environmentSlug + ); + + await keyStore.setItemWithExpiry( + key, + KeyStoreTtls.SendFailedIntegrationSyncEmailsInSeconds, + JSON.stringify(payload) + ); + await queueService.queue( + QueueName.IntegrationSync, + QueueJobs.SendFailedIntegrationSyncEmails, + { + secretPath: payload.secretPath, + projectId: payload.projectId, + environmentSlug: payload.environmentSlug + }, + { + delay: applyJitter( + secondsToMillis(KeyStoreTtls.SendFailedIntegrationSyncEmailsInSeconds / 2), + secondsToMillis(10) + ), + jobId: key, + removeOnFail: true, + removeOnComplete: true + } + ); + }; + queueService.start(QueueName.SecretSync, async (job) => { const { _deDupeQueue: deDupeQueue, @@ -560,371 +595,416 @@ export const secretQueueFactory = ({ }); queueService.start(QueueName.IntegrationSync, async (job) => { - const { environment, actorId, isManual, projectId, secretPath, depth = 1, deDupeQueue = {} } = job.data; - if (depth > MAX_SYNC_SECRET_DEPTH) return; - - const folder = await folderDAL.findBySecretPath(projectId, environment, secretPath); - if (!folder) { - throw new Error("Secret path not found"); - } - - const sendIntegrationSyncFailedMail = async (integrations: { integrationId: string; syncMessage?: string }[]) => { + if (job.name === QueueJobs.SendFailedIntegrationSyncEmails) { const appCfg = getConfig(); - // If smtp is not configured, we can return early without having to fetch the project members. + // If smtp is not configured, we can return early if (!appCfg.isSmtpConfigured) return; - const projectMembers = await projectMembershipDAL.findAllProjectMembers(projectId); - const project = await projectDAL.findById(projectId); + const jobPayload = job.data as Pick< + TFailedIntegrationSyncEmailsPayload, + "projectId" | "secretPath" | "environmentSlug" + >; + + const failedIntegrationsDetails = await keyStore.getItem( + KeyStorePrefixes.SendFailedIntegrationSyncEmails( + jobPayload.projectId, + jobPayload.secretPath, + jobPayload.environmentSlug + ) + ); + + if (!failedIntegrationsDetails) return; + + const failedSyncKeyStore = FailedIntegrationSyncEmailsPayloadSchema.parse(JSON.parse(failedIntegrationsDetails)); + + const projectMembers = await projectMembershipDAL.findAllProjectMembers(failedSyncKeyStore.projectId); + const project = await projectDAL.findById(failedSyncKeyStore.projectId); // Only send emails to admins, and if its a manual trigger, only send it to the person who triggered it (if actor is admin as well) const filteredProjectMembers = projectMembers .filter((member) => member.roles.some((role) => role.role === ProjectMembershipRole.Admin)) - .filter((member) => (isManual && actorId ? member.userId === actorId : true)); + .filter((member) => + failedSyncKeyStore.manuallyTriggeredByUserId + ? member.userId === failedSyncKeyStore.manuallyTriggeredByUserId + : true + ); await smtpService.sendMail({ recipients: filteredProjectMembers.map((member) => member.user.email!), template: SmtpTemplates.IntegrationSyncFailed, subjectLine: `Integration Sync Failed`, substitutions: { - syncMessage: integrations[0]?.syncMessage, // We are only displaying the sync message if its a singular integration, so we can just grab the first one in the array. - secretPath, - environment: folder.environment.name, - count: integrations.length, + syncMessage: failedSyncKeyStore.count === 1 ? failedSyncKeyStore.syncMessage : undefined, // We are only displaying the sync message if its a singular integration, so we can just grab the first one in the array. + secretPath: failedSyncKeyStore.secretPath, + environment: failedSyncKeyStore.environmentName, + count: failedSyncKeyStore.count, projectName: project.name, integrationUrl: `${appCfg.SITE_URL}/integrations/${project.id}` } }); - }; - - // find all imports made with the given environment and secret path - const linkSourceDto = { - projectId, - importEnv: folder.environment.id, - importPath: secretPath, - isReplication: false - }; - const imports = await secretImportDAL.find(linkSourceDto); - - if (imports.length) { - // 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); - logger.info( - `getIntegrationSecrets: Syncing secret due to link change [jobId=${job.id}] [projectId=${job.data.projectId}] [environment=${job.data.environment}] [secretPath=${job.data.secretPath}] [depth=${depth}]` - ); - await Promise.all( - imports - .filter(({ folderId }) => Boolean(foldersGroupedById[folderId][0]?.path as string)) - // filter out already synced ones - .filter( - ({ folderId }) => - !deDupeQueue[ - uniqueSecretQueueKey( - foldersGroupedById[folderId][0]?.environmentSlug as string, - foldersGroupedById[folderId][0]?.path as string - ) - ] - ) - .map(({ folderId }) => - syncSecrets({ - projectId, - secretPath: foldersGroupedById[folderId][0]?.path as string, - environmentSlug: foldersGroupedById[folderId][0]?.environmentSlug as string, - _deDupeQueue: deDupeQueue, - _depth: depth + 1, - excludeReplication: true - }) - ) - ); } - const { shouldUseSecretV2Bridge, botKey } = await projectBotService.getBotKey(projectId); - const { decryptor: secretManagerDecryptor } = await kmsService.createCipherPairWithDataKey({ - type: KmsDataKey.SecretManager, - projectId - }); - let referencedFolderIds; - if (shouldUseSecretV2Bridge) { - const secretReferences = await secretV2BridgeDAL.findReferencedSecretReferences( + + if (job.name === QueueJobs.IntegrationSync) { + const { + environment, + actorId, + isManual, projectId, - folder.environment.slug, - secretPath - ); - referencedFolderIds = unique(secretReferences, (i) => i.folderId).map(({ folderId }) => folderId); - } else { - const secretReferences = await secretDAL.findReferencedSecretReferences( - projectId, - folder.environment.slug, - secretPath - ); - referencedFolderIds = unique(secretReferences, (i) => i.folderId).map(({ folderId }) => folderId); - } - if (referencedFolderIds.length) { - const referencedFolders = await folderDAL.findSecretPathByFolderIds(projectId, referencedFolderIds); - const referencedFoldersGroupedById = groupBy(referencedFolders.filter(Boolean), (i) => i?.id as string); - logger.info( - `getIntegrationSecrets: Syncing secret due to reference change [jobId=${job.id}] [projectId=${job.data.projectId}] [environment=${job.data.environment}] [secretPath=${job.data.secretPath}] [depth=${depth}]` - ); - await Promise.all( - referencedFolderIds - .filter((folderId) => Boolean(referencedFoldersGroupedById[folderId][0]?.path)) - // filter out already synced ones - .filter( - (folderId) => - !deDupeQueue[ - uniqueSecretQueueKey( - referencedFoldersGroupedById[folderId][0]?.environmentSlug as string, - referencedFoldersGroupedById[folderId][0]?.path as string - ) - ] - ) - .map((folderId) => - syncSecrets({ - projectId, - secretPath: referencedFoldersGroupedById[folderId][0]?.path as string, - environmentSlug: referencedFoldersGroupedById[folderId][0]?.environmentSlug as string, - _deDupeQueue: deDupeQueue, - _depth: depth + 1, - excludeReplication: true - }) - ) - ); - } + secretPath, + depth = 1, + deDupeQueue = {} + } = job.data as TIntegrationSyncPayload; + if (depth > MAX_SYNC_SECRET_DEPTH) return; - const integrations = await integrationDAL.findByProjectIdV2(projectId, environment); // note: returns array of integrations + integration auths in this environment - const toBeSyncedIntegrations = integrations.filter( - // note: sync only the integrations sourced from secretPath - ({ secretPath: integrationSecPath, isActive }) => isActive && isSamePath(secretPath, integrationSecPath) - ); - - const integrationsFailedToSync: { integrationId: string }[] = []; - - if (!integrations.length) return; - logger.info( - `getIntegrationSecrets: secret integration sync started [jobId=${job.id}] [jobId=${job.id}] [projectId=${job.data.projectId}] [environment=${job.data.environment}] [secretPath=${job.data.secretPath}] [depth=${job.data.depth}]` - ); - - const lock = await keyStore.acquireLock( - [KeyStorePrefixes.SyncSecretIntegrationLock(projectId, environment, secretPath)], - 10000, - { - retryCount: 3, - retryDelay: 2000 + const folder = await folderDAL.findBySecretPath(projectId, environment, secretPath); + if (!folder) { + throw new Error("Secret path not found"); } - ); - const lockAcquiredTime = new Date(); - const lastRunSyncIntegrationTimestamp = await keyStore.getItem( - KeyStorePrefixes.SyncSecretIntegrationLastRunTimestamp(projectId, environment, secretPath) - ); + // find all imports made with the given environment and secret path + const linkSourceDto = { + projectId, + importEnv: folder.environment.id, + importPath: secretPath, + isReplication: false + }; + const imports = await secretImportDAL.find(linkSourceDto); - // check whether the integration should wait or not - if (lastRunSyncIntegrationTimestamp) { - const INTEGRATION_INTERVAL = 2000; - const isStaleSyncIntegration = new Date(job.timestamp) < new Date(lastRunSyncIntegrationTimestamp); - if (isStaleSyncIntegration) { + if (imports.length) { + // 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); logger.info( - `getIntegrationSecrets: secret integration sync stale [jobId=${job.id}] [jobId=${job.id}] [projectId=${job.data.projectId}] [environment=${job.data.environment}] [secretPath=${job.data.secretPath}] [depth=${job.data.depth}]` + `getIntegrationSecrets: Syncing secret due to link change [jobId=${job.id}] [projectId=${job.data.projectId}] [environment=${environment}] [secretPath=${job.data.secretPath}] [depth=${depth}]` + ); + await Promise.all( + imports + .filter(({ folderId }) => Boolean(foldersGroupedById[folderId][0]?.path as string)) + // filter out already synced ones + .filter( + ({ folderId }) => + !deDupeQueue[ + uniqueSecretQueueKey( + foldersGroupedById[folderId][0]?.environmentSlug as string, + foldersGroupedById[folderId][0]?.path as string + ) + ] + ) + .map(({ folderId }) => + syncSecrets({ + projectId, + secretPath: foldersGroupedById[folderId][0]?.path as string, + environmentSlug: foldersGroupedById[folderId][0]?.environmentSlug as string, + _deDupeQueue: deDupeQueue, + _depth: depth + 1, + excludeReplication: true + }) + ) + ); + } + const { shouldUseSecretV2Bridge, botKey } = await projectBotService.getBotKey(projectId); + const { decryptor: secretManagerDecryptor } = await kmsService.createCipherPairWithDataKey({ + type: KmsDataKey.SecretManager, + projectId + }); + let referencedFolderIds; + if (shouldUseSecretV2Bridge) { + const secretReferences = await secretV2BridgeDAL.findReferencedSecretReferences( + projectId, + folder.environment.slug, + secretPath + ); + referencedFolderIds = unique(secretReferences, (i) => i.folderId).map(({ folderId }) => folderId); + } else { + const secretReferences = await secretDAL.findReferencedSecretReferences( + projectId, + folder.environment.slug, + secretPath + ); + referencedFolderIds = unique(secretReferences, (i) => i.folderId).map(({ folderId }) => folderId); + } + if (referencedFolderIds.length) { + const referencedFolders = await folderDAL.findSecretPathByFolderIds(projectId, referencedFolderIds); + const referencedFoldersGroupedById = groupBy(referencedFolders.filter(Boolean), (i) => i?.id as string); + logger.info( + `getIntegrationSecrets: Syncing secret due to reference change [jobId=${job.id}] [projectId=${job.data.projectId}] [environment=${environment}] [secretPath=${job.data.secretPath}] [depth=${depth}]` + ); + await Promise.all( + referencedFolderIds + .filter((folderId) => Boolean(referencedFoldersGroupedById[folderId][0]?.path)) + // filter out already synced ones + .filter( + (folderId) => + !deDupeQueue[ + uniqueSecretQueueKey( + referencedFoldersGroupedById[folderId][0]?.environmentSlug as string, + referencedFoldersGroupedById[folderId][0]?.path as string + ) + ] + ) + .map((folderId) => + syncSecrets({ + projectId, + secretPath: referencedFoldersGroupedById[folderId][0]?.path as string, + environmentSlug: referencedFoldersGroupedById[folderId][0]?.environmentSlug as string, + _deDupeQueue: deDupeQueue, + _depth: depth + 1, + excludeReplication: true + }) + ) ); - return; } - const timeDifferenceWithLastIntegration = getTimeDifferenceInSeconds( - lockAcquiredTime.toISOString(), - lastRunSyncIntegrationTimestamp + const integrations = await integrationDAL.findByProjectIdV2(projectId, environment); // note: returns array of integrations + integration auths in this environment + const toBeSyncedIntegrations = integrations.filter( + // note: sync only the integrations sourced from secretPath + ({ secretPath: integrationSecPath, isActive }) => isActive && isSamePath(secretPath, integrationSecPath) ); - if (timeDifferenceWithLastIntegration < INTEGRATION_INTERVAL && timeDifferenceWithLastIntegration > 0) - await new Promise((resolve) => { - setTimeout(resolve, 2000 - timeDifferenceWithLastIntegration * 1000); - }); - } - const generateActor = async (): Promise => { - if (isManual && actorId) { - const user = await userDAL.findById(actorId); + const integrationsFailedToSync: { integrationId: string; syncMessage?: string }[] = []; - if (!user) { - throw new Error("User not found"); + if (!integrations.length) return; + logger.info( + `getIntegrationSecrets: secret integration sync started [jobId=${job.id}] [jobId=${job.id}] [projectId=${job.data.projectId}] [environment=${environment}] [secretPath=${job.data.secretPath}] [depth=${depth}]` + ); + + const lock = await keyStore.acquireLock( + [KeyStorePrefixes.SyncSecretIntegrationLock(projectId, environment, secretPath)], + 10000, + { + retryCount: 3, + retryDelay: 2000 + } + ); + const lockAcquiredTime = new Date(); + + const lastRunSyncIntegrationTimestamp = await keyStore.getItem( + KeyStorePrefixes.SyncSecretIntegrationLastRunTimestamp(projectId, environment, secretPath) + ); + + // check whether the integration should wait or not + if (lastRunSyncIntegrationTimestamp) { + const INTEGRATION_INTERVAL = 2000; + const isStaleSyncIntegration = new Date(job.timestamp) < new Date(lastRunSyncIntegrationTimestamp); + if (isStaleSyncIntegration) { + logger.info( + `getIntegrationSecrets: secret integration sync stale [jobId=${job.id}] [jobId=${job.id}] [projectId=${job.data.projectId}] [environment=${environment}] [secretPath=${job.data.secretPath}] [depth=${depth}]` + ); + return; + } + + const timeDifferenceWithLastIntegration = getTimeDifferenceInSeconds( + lockAcquiredTime.toISOString(), + lastRunSyncIntegrationTimestamp + ); + if (timeDifferenceWithLastIntegration < INTEGRATION_INTERVAL && timeDifferenceWithLastIntegration > 0) + await new Promise((resolve) => { + setTimeout(resolve, 2000 - timeDifferenceWithLastIntegration * 1000); + }); + } + + const generateActor = async (): Promise => { + if (isManual && actorId) { + const user = await userDAL.findById(actorId); + + if (!user) { + throw new Error("User not found"); + } + + return { + type: ActorType.USER, + metadata: { + email: user.email, + username: user.username, + userId: user.id + } + }; } return { - type: ActorType.USER, - metadata: { - email: user.email, - username: user.username, - userId: user.id - } + type: ActorType.PLATFORM, + metadata: {} }; - } - - return { - type: ActorType.PLATFORM, - metadata: {} }; - }; - // akhilmhdh: this try catch is for lock release - try { - const secrets = shouldUseSecretV2Bridge - ? await getIntegrationSecretsV2({ - environment, - projectId, - folderId: folder.id, - depth: 1, - secretPath, - decryptor: (value) => (value ? secretManagerDecryptor({ cipherTextBlob: value }).toString() : "") - }) - : await getIntegrationSecrets({ - environment, - projectId, - folderId: folder.id, - key: botKey as string, - depth: 1, - secretPath - }); + // akhilmhdh: this try catch is for lock release + try { + const secrets = shouldUseSecretV2Bridge + ? await getIntegrationSecretsV2({ + environment, + projectId, + folderId: folder.id, + depth: 1, + secretPath, + decryptor: (value) => (value ? secretManagerDecryptor({ cipherTextBlob: value }).toString() : "") + }) + : await getIntegrationSecrets({ + environment, + projectId, + folderId: folder.id, + key: botKey as string, + depth: 1, + secretPath + }); - for (const integration of toBeSyncedIntegrations) { - const integrationAuth = { - ...integration.integrationAuth, - createdAt: new Date(), - updatedAt: new Date(), - projectId: integration.projectId - }; + for (const integration of toBeSyncedIntegrations) { + const integrationAuth = { + ...integration.integrationAuth, + createdAt: new Date(), + updatedAt: new Date(), + projectId: integration.projectId + }; - const { accessToken, accessId } = await integrationAuthService.getIntegrationAccessToken( - integrationAuth, - shouldUseSecretV2Bridge, - botKey - ); - let awsAssumeRoleArn = null; - if (shouldUseSecretV2Bridge) { - if (integrationAuth.encryptedAwsAssumeIamRoleArn) { - awsAssumeRoleArn = secretManagerDecryptor({ - cipherTextBlob: Buffer.from(integrationAuth.encryptedAwsAssumeIamRoleArn) - }).toString(); - } - } else if ( - integrationAuth.awsAssumeIamRoleArnTag && - integrationAuth.awsAssumeIamRoleArnIV && - integrationAuth.awsAssumeIamRoleArnCipherText - ) { - awsAssumeRoleArn = decryptSymmetric128BitHexKeyUTF8({ - ciphertext: integrationAuth.awsAssumeIamRoleArnCipherText, - iv: integrationAuth.awsAssumeIamRoleArnIV, - tag: integrationAuth.awsAssumeIamRoleArnTag, - key: botKey as string - }); - } - - const suffixedSecrets: typeof secrets = {}; - const metadata = integration.metadata as Record; - if (metadata) { - Object.keys(secrets).forEach((key) => { - const prefix = metadata?.secretPrefix || ""; - const suffix = metadata?.secretSuffix || ""; - const newKey = prefix + key + suffix; - suffixedSecrets[newKey] = secrets[key]; - }); - } - - // akhilmhdh: this try catch is for catching integration error and saving it in db - try { - // akhilmhdh: this needs to changed later to be more easier to use - // at present this is not at all extendable like to add a new parameter for just one integration need to modify multiple places - const response = await syncIntegrationSecrets({ - createManySecretsRawFn, - updateManySecretsRawFn, - integrationDAL, - integration, + const { accessToken, accessId } = await integrationAuthService.getIntegrationAccessToken( integrationAuth, - secrets: Object.keys(suffixedSecrets).length !== 0 ? suffixedSecrets : secrets, - accessId: accessId as string, - awsAssumeRoleArn, - accessToken, - projectId, - appendices: { - prefix: metadata?.secretPrefix || "", - suffix: metadata?.secretSuffix || "" + shouldUseSecretV2Bridge, + botKey + ); + let awsAssumeRoleArn = null; + if (shouldUseSecretV2Bridge) { + if (integrationAuth.encryptedAwsAssumeIamRoleArn) { + awsAssumeRoleArn = secretManagerDecryptor({ + cipherTextBlob: Buffer.from(integrationAuth.encryptedAwsAssumeIamRoleArn) + }).toString(); } - }); - - await auditLogService.createAuditLog({ - projectId, - actor: await generateActor(), - event: { - type: EventType.INTEGRATION_SYNCED, - metadata: { - integrationId: integration.id, - isSynced: response?.isSynced ?? true, - lastSyncJobId: job?.id ?? "", - lastUsed: new Date(), - syncMessage: response?.syncMessage ?? "" - } - } - }); - - await integrationDAL.updateById(integration.id, { - lastSyncJobId: job.id, - lastUsed: new Date(), - syncMessage: response?.syncMessage ?? "", - isSynced: response?.isSynced ?? true - }); - - // May be undefined, if it's undefined we assume the sync was successful, hence the strict equality type check. - if (response?.isSynced === false) { - integrationsFailedToSync.push({ - integrationId: integration.id + } else if ( + integrationAuth.awsAssumeIamRoleArnTag && + integrationAuth.awsAssumeIamRoleArnIV && + integrationAuth.awsAssumeIamRoleArnCipherText + ) { + awsAssumeRoleArn = decryptSymmetric128BitHexKeyUTF8({ + ciphertext: integrationAuth.awsAssumeIamRoleArnCipherText, + iv: integrationAuth.awsAssumeIamRoleArnIV, + tag: integrationAuth.awsAssumeIamRoleArnTag, + key: botKey as string }); } - } catch (err) { - logger.error( - err, - `Secret integration sync error [projectId=${job.data.projectId}] [environment=${job.data.environment}] [secretPath=${job.data.secretPath}]` - ); - const message = - (err instanceof AxiosError ? JSON.stringify(err?.response?.data) : (err as Error)?.message) || - "Unknown error occurred."; + const suffixedSecrets: typeof secrets = {}; + const metadata = integration.metadata as Record; + if (metadata) { + Object.keys(secrets).forEach((key) => { + const prefix = metadata?.secretPrefix || ""; + const suffix = metadata?.secretSuffix || ""; + const newKey = prefix + key + suffix; + suffixedSecrets[newKey] = secrets[key]; + }); + } - await auditLogService.createAuditLog({ - projectId, - actor: await generateActor(), - event: { - type: EventType.INTEGRATION_SYNCED, - metadata: { - integrationId: integration.id, - isSynced: false, - lastSyncJobId: job?.id ?? "", - lastUsed: new Date(), - syncMessage: message + // akhilmhdh: this try catch is for catching integration error and saving it in db + try { + // akhilmhdh: this needs to changed later to be more easier to use + // at present this is not at all extendable like to add a new parameter for just one integration need to modify multiple places + const response = await syncIntegrationSecrets({ + createManySecretsRawFn, + updateManySecretsRawFn, + integrationDAL, + integration, + integrationAuth, + secrets: Object.keys(suffixedSecrets).length !== 0 ? suffixedSecrets : secrets, + accessId: accessId as string, + awsAssumeRoleArn, + accessToken, + projectId, + appendices: { + prefix: metadata?.secretPrefix || "", + suffix: metadata?.secretSuffix || "" } + }); + + await auditLogService.createAuditLog({ + projectId, + actor: await generateActor(), + event: { + type: EventType.INTEGRATION_SYNCED, + metadata: { + integrationId: integration.id, + isSynced: response?.isSynced ?? true, + lastSyncJobId: job?.id ?? "", + lastUsed: new Date(), + syncMessage: response?.syncMessage ?? "" + } + } + }); + + await integrationDAL.updateById(integration.id, { + lastSyncJobId: job.id, + lastUsed: new Date(), + syncMessage: response?.syncMessage ?? "", + isSynced: response?.isSynced ?? true + }); + + // May be undefined, if it's undefined we assume the sync was successful, hence the strict equality type check. + if (response?.isSynced === false) { + integrationsFailedToSync.push({ + integrationId: integration.id, + syncMessage: response.syncMessage + }); } - }); + } catch (err) { + logger.error( + err, + `Secret integration sync error [projectId=${job.data.projectId}] [environment=${environment}] [secretPath=${job.data.secretPath}]` + ); - await integrationDAL.updateById(integration.id, { - lastSyncJobId: job.id, - syncMessage: message, - isSynced: false - }); + const message = + (err instanceof AxiosError ? JSON.stringify(err?.response?.data) : (err as Error)?.message) || + "Unknown error occurred."; - integrationsFailedToSync.push({ - integrationId: integration.id + await auditLogService.createAuditLog({ + projectId, + actor: await generateActor(), + event: { + type: EventType.INTEGRATION_SYNCED, + metadata: { + integrationId: integration.id, + isSynced: false, + lastSyncJobId: job?.id ?? "", + lastUsed: new Date(), + syncMessage: message + } + } + }); + + await integrationDAL.updateById(integration.id, { + lastSyncJobId: job.id, + syncMessage: message, + isSynced: false + }); + + integrationsFailedToSync.push({ + integrationId: integration.id, + syncMessage: message + }); + } + } + } finally { + await lock.release(); + if (integrationsFailedToSync.length) { + await sendFailedIntegrationSyncEmails({ + count: integrationsFailedToSync.length, + environmentName: folder.environment.name, + environmentSlug: environment, + ...(isManual && + actorId && { + manuallyTriggeredByUserId: actorId + }), + projectId, + secretPath, + syncMessage: integrationsFailedToSync[0].syncMessage }); } } - } finally { - await lock.release(); - await sendIntegrationSyncFailedMail(integrationsFailedToSync); + await keyStore.setItemWithExpiry( + KeyStorePrefixes.SyncSecretIntegrationLastRunTimestamp(projectId, environment, secretPath), + KeyStoreTtls.SetSyncSecretIntegrationLastRunTimestampInSeconds, + lockAcquiredTime.toISOString() + ); + logger.info("Secret integration sync ended: %s", job.id); } - - await keyStore.setItemWithExpiry( - KeyStorePrefixes.SyncSecretIntegrationLastRunTimestamp(projectId, environment, secretPath), - KeyStoreTtls.SetSyncSecretIntegrationLastRunTimestampInSeconds, - lockAcquiredTime.toISOString() - ); - logger.info("Secret integration sync ended: %s", job.id); }); queueService.start(QueueName.SecretReminder, async ({ data }) => { diff --git a/backend/src/services/secret/secret-types.ts b/backend/src/services/secret/secret-types.ts index 1686ee488..b5281af68 100644 --- a/backend/src/services/secret/secret-types.ts +++ b/backend/src/services/secret/secret-types.ts @@ -1,4 +1,5 @@ import { Knex } from "knex"; +import { z } from "zod"; import { SecretType, TSecretBlindIndexes, TSecrets, TSecretsInsert, TSecretsUpdate } from "@app/db/schemas"; import { TProjectPermission } from "@app/lib/types"; @@ -21,6 +22,29 @@ type TPartialSecret = Pick; +export const FailedIntegrationSyncEmailsPayloadSchema = z.object({ + projectId: z.string(), + secretPath: z.string(), + environmentName: z.string(), + environmentSlug: z.string(), + + count: z.number(), + syncMessage: z.string().optional(), + manuallyTriggeredByUserId: z.string().optional() +}); + +export type TFailedIntegrationSyncEmailsPayload = z.infer; + +export type TIntegrationSyncPayload = { + isManual?: boolean; + actorId?: string; + projectId: string; + environment: string; + secretPath: string; + depth?: number; + deDupeQueue?: Record; +}; + export type TCreateSecretDTO = { secretName: string; path: string;