Merge pull request #1988 from Infisical/daniel/ingrations-improvements

Fix: Silent integration errors
This commit is contained in:
Daniel Hougaard
2024-06-18 15:14:23 +02:00
committed by GitHub
2 changed files with 112 additions and 101 deletions
@@ -18,7 +18,7 @@ import {
UpdateSecretCommand UpdateSecretCommand
} from "@aws-sdk/client-secrets-manager"; } from "@aws-sdk/client-secrets-manager";
import { Octokit } from "@octokit/rest"; import { Octokit } from "@octokit/rest";
import AWS from "aws-sdk"; import AWS, { AWSError } from "aws-sdk";
import { AxiosError } from "axios"; import { AxiosError } from "axios";
import sodium from "libsodium-wrappers"; import sodium from "libsodium-wrappers";
import isEqual from "lodash.isequal"; import isEqual from "lodash.isequal";
@@ -452,7 +452,11 @@ const syncSecretsAWSParameterStore = async ({
accessId: string | null; accessId: string | null;
accessToken: string; accessToken: string;
}) => { }) => {
if (!accessId) return; let response: { isSynced: boolean; syncMessage: string } | null = null;
if (!accessId) {
throw new Error("AWS access ID is required");
}
const config = new AWS.Config({ const config = new AWS.Config({
region: integration.region as string, region: integration.region as string,
@@ -557,6 +561,11 @@ const syncSecretsAWSParameterStore = async ({
`AWS Parameter Store Error [integration=${integration.id}]: double check AWS account permissions (refer to the Infisical docs)` `AWS Parameter Store Error [integration=${integration.id}]: double check AWS account permissions (refer to the Infisical docs)`
); );
} }
response = {
isSynced: false,
syncMessage: (err as AWSError)?.message || "Error syncing with AWS Parameter Store"
};
} }
} }
} }
@@ -585,6 +594,8 @@ const syncSecretsAWSParameterStore = async ({
} }
} }
} }
return response;
}; };
/** /**
@@ -603,7 +614,9 @@ const syncSecretsAWSSecretManager = async ({
}) => { }) => {
const metadata = z.record(z.any()).parse(integration.metadata || {}); const metadata = z.record(z.any()).parse(integration.metadata || {});
if (!accessId) return; if (!accessId) {
throw new Error("AWS access ID is required");
}
const secretsManager = new SecretsManagerClient({ const secretsManager = new SecretsManagerClient({
region: integration.region as string, region: integration.region as string,
@@ -722,7 +735,7 @@ const syncSecretsAWSSecretManager = async ({
} }
} }
} catch (err) { } catch (err) {
// case when AWS manager can't find the specified secret // case 1: when AWS manager can't find the specified secret
if (err instanceof ResourceNotFoundException && secretsManager) { if (err instanceof ResourceNotFoundException && secretsManager) {
await secretsManager.send( await secretsManager.send(
new CreateSecretCommand({ new CreateSecretCommand({
@@ -734,6 +747,9 @@ const syncSecretsAWSSecretManager = async ({
: [] : []
}) })
); );
// case 2: something unexpected went wrong, so we'll throw the error to reflect the error in the integration sync status
} else {
throw err;
} }
} }
}; };
@@ -753,14 +769,12 @@ const syncSecretsAWSSecretManager = async ({
const syncSecretsHeroku = async ({ const syncSecretsHeroku = async ({
createManySecretsRawFn, createManySecretsRawFn,
updateManySecretsRawFn, updateManySecretsRawFn,
integrationDAL,
integration, integration,
secrets, secrets,
accessToken accessToken
}: { }: {
createManySecretsRawFn: (params: TCreateManySecretsRawFn) => Promise<Array<TSecrets & { _id: string }>>; createManySecretsRawFn: (params: TCreateManySecretsRawFn) => Promise<Array<TSecrets & { _id: string }>>;
updateManySecretsRawFn: (params: TUpdateManySecretsRawFn) => Promise<Array<TSecrets & { _id: string }>>; updateManySecretsRawFn: (params: TUpdateManySecretsRawFn) => Promise<Array<TSecrets & { _id: string }>>;
integrationDAL: Pick<TIntegrationDALFactory, "updateById">;
integration: TIntegrations & { integration: TIntegrations & {
projectId: string; projectId: string;
environment: { environment: {
@@ -862,10 +876,6 @@ const syncSecretsHeroku = async ({
} }
} }
); );
await integrationDAL.updateById(integration.id, {
lastUsed: new Date()
});
}; };
/** /**
@@ -2656,7 +2666,9 @@ const syncSecretsHashiCorpVault = async ({
accessId: string | null; accessId: string | null;
accessToken: string; accessToken: string;
}) => { }) => {
if (!accessId) return; if (!accessId) {
throw new Error("Access ID is required");
}
interface LoginAppRoleRes { interface LoginAppRoleRes {
auth: { auth: {
@@ -3486,6 +3498,8 @@ export const syncIntegrationSecrets = async ({
accessToken: string; accessToken: string;
appendices?: { prefix: string; suffix: string }; appendices?: { prefix: string; suffix: string };
}) => { }) => {
let response: { isSynced: boolean; syncMessage: string } | null = null;
switch (integration.integration) { switch (integration.integration) {
case Integrations.GCP_SECRET_MANAGER: case Integrations.GCP_SECRET_MANAGER:
await syncSecretsGCPSecretManager({ await syncSecretsGCPSecretManager({
@@ -3502,7 +3516,7 @@ export const syncIntegrationSecrets = async ({
}); });
break; break;
case Integrations.AWS_PARAMETER_STORE: case Integrations.AWS_PARAMETER_STORE:
await syncSecretsAWSParameterStore({ response = await syncSecretsAWSParameterStore({
integration, integration,
secrets, secrets,
accessId, accessId,
@@ -3521,7 +3535,6 @@ export const syncIntegrationSecrets = async ({
await syncSecretsHeroku({ await syncSecretsHeroku({
createManySecretsRawFn, createManySecretsRawFn,
updateManySecretsRawFn, updateManySecretsRawFn,
integrationDAL,
integration, integration,
secrets, secrets,
accessToken accessToken
@@ -3727,4 +3740,6 @@ export const syncIntegrationSecrets = async ({
default: default:
throw new BadRequestError({ message: "Invalid integration" }); throw new BadRequestError({ message: "Invalid integration" });
} }
return response;
}; };
+84 -88
View File
@@ -421,94 +421,88 @@ export const secretQueueFactory = ({
const folder = await folderDAL.findBySecretPath(projectId, environment, secretPath); const folder = await folderDAL.findBySecretPath(projectId, environment, secretPath);
if (!folder) { if (!folder) {
logger.error(new Error("Secret path not found")); throw new Error("Secret path not found");
return;
} }
// start syncing all linked imports also // find all imports made with the given environment and secret path
if (depth < MAX_SYNC_SECRET_DEPTH) { const linkSourceDto = {
// find all imports made with the given environment and secret path projectId,
const linkSourceDto = { importEnv: folder.environment.id,
projectId, importPath: secretPath,
importEnv: folder.environment.id, isReplication: false
importPath: secretPath, };
isReplication: false 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);
logger.info( 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}]` `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( await Promise.all(
imports imports
.filter(({ folderId }) => Boolean(foldersGroupedById[folderId][0]?.path as string)) .filter(({ folderId }) => Boolean(foldersGroupedById[folderId][0]?.path as string))
// filter out already synced ones // filter out already synced ones
.filter( .filter(
({ folderId }) => ({ folderId }) =>
!deDupeQueue[ !deDupeQueue[
uniqueSecretQueueKey( uniqueSecretQueueKey(
foldersGroupedById[folderId][0]?.environmentSlug as string, foldersGroupedById[folderId][0]?.environmentSlug as string,
foldersGroupedById[folderId][0]?.path as string foldersGroupedById[folderId][0]?.path as string
) )
] ]
) )
.map(({ folderId }) => .map(({ folderId }) =>
syncSecrets({ syncSecrets({
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,
_deDupeQueue: deDupeQueue, _deDupeQueue: deDupeQueue,
_depth: depth + 1, _depth: depth + 1,
excludeReplication: true excludeReplication: true
}) })
) )
); );
} }
const secretReferences = await secretDAL.findReferencedSecretReferences( const secretReferences = await secretDAL.findReferencedSecretReferences(
projectId, projectId,
folder.environment.slug, folder.environment.slug,
secretPath secretPath
);
if (secretReferences.length) {
const referencedFolderIds = unique(secretReferences, (i) => i.folderId).map(({ folderId }) => folderId);
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(
secretReferences
.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
})
)
); );
if (secretReferences.length) {
const referencedFolderIds = unique(secretReferences, (i) => i.folderId).map(({ folderId }) => folderId);
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(
secretReferences
.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
})
)
);
}
} else {
logger.info(`getIntegrationSecrets: Secret depth exceeded for [projectId=${projectId}] [folderId=${folder.id}]`);
} }
const integrations = await integrationDAL.findByProjectIdV2(projectId, environment); // note: returns array of integrations + integration auths in this environment const integrations = await integrationDAL.findByProjectIdV2(projectId, environment); // note: returns array of integrations + integration auths in this environment
@@ -550,7 +544,7 @@ export const secretQueueFactory = ({
} }
try { try {
await syncIntegrationSecrets({ const response = await syncIntegrationSecrets({
createManySecretsRawFn, createManySecretsRawFn,
updateManySecretsRawFn, updateManySecretsRawFn,
integrationDAL, integrationDAL,
@@ -568,13 +562,15 @@ export const secretQueueFactory = ({
await integrationDAL.updateById(integration.id, { await integrationDAL.updateById(integration.id, {
lastSyncJobId: job.id, lastSyncJobId: job.id,
lastUsed: new Date(), lastUsed: new Date(),
syncMessage: "", syncMessage: response?.syncMessage ?? "",
isSynced: true isSynced: response?.isSynced ?? true
}); });
} catch (err: unknown) { } catch (err) {
logger.info("Secret integration sync error: %o", err); logger.info("Secret integration sync error: %o", err);
const message = const message =
err instanceof AxiosError ? JSON.stringify((err as AxiosError)?.response?.data) : (err as Error)?.message; (err instanceof AxiosError ? JSON.stringify(err?.response?.data) : (err as Error)?.message) ||
"Unknown error occurred.";
await integrationDAL.updateById(integration.id, { await integrationDAL.updateById(integration.id, {
lastSyncJobId: job.id, lastSyncJobId: job.id,