Address PR comments

This commit is contained in:
Carlos Monastyrski
2025-09-19 03:43:22 -03:00
parent 8130be5e2f
commit 27dbe5b013
72 changed files with 2913 additions and 1367 deletions

View File

@@ -19,7 +19,7 @@ export async function up(knex: Knex): Promise<void> {
t.uuid("subscriberId");
t.foreign("subscriberId").references("id").inTable(TableName.PkiSubscriber).onDelete("SET NULL");
t.uuid("connectionId").notNullable();
t.foreign("connectionId").references("id").inTable(TableName.AppConnection).onDelete("CASCADE");
t.foreign("connectionId").references("id").inTable(TableName.AppConnection);
t.timestamps(true, true, true);
t.string("syncStatus");
t.string("lastSyncJobId");

View File

@@ -256,7 +256,7 @@ export type SecretSyncSubjectFields = {
};
export type PkiSyncSubjectFields = {
projectId: string;
subscriberName: string;
};
export type DynamicSecretSubjectFields = {
@@ -501,7 +501,17 @@ const SecretSyncConditionV2Schema = z
const PkiSyncConditionSchema = z
.object({
projectId: z.string()
subscriberName: z.union([
z.string(),
z
.object({
[PermissionConditionOperators.$EQ]: PermissionConditionSchema[PermissionConditionOperators.$EQ],
[PermissionConditionOperators.$NEQ]: PermissionConditionSchema[PermissionConditionOperators.$NEQ],
[PermissionConditionOperators.$IN]: PermissionConditionSchema[PermissionConditionOperators.$IN],
[PermissionConditionOperators.$GLOB]: PermissionConditionSchema[PermissionConditionOperators.$GLOB]
})
.partial()
])
})
.partial();

View File

@@ -27,8 +27,7 @@ import { TCreateUserNotificationDTO } from "@app/services/notification/notificat
import {
TQueuePkiSyncImportCertificatesByIdDTO,
TQueuePkiSyncRemoveCertificatesByIdDTO,
TQueuePkiSyncSyncCertificatesByIdDTO,
TQueueSendPkiSyncActionFailedNotificationsDTO
TQueuePkiSyncSyncCertificatesByIdDTO
} from "@app/services/pki-sync/pki-sync-types";
import {
TFailedIntegrationSyncEmailsPayload,
@@ -110,7 +109,6 @@ export enum QueueJobs {
PkiSyncSyncCertificates = "pki-sync-sync-certificates",
PkiSyncImportCertificates = "pki-sync-import-certificates",
PkiSyncRemoveCertificates = "pki-sync-remove-certificates",
PkiSyncSendActionFailedNotifications = "pki-sync-send-action-failed-notifications",
SecretRotationV2QueueRotations = "secret-rotation-v2-queue-rotations",
SecretRotationV2RotateSecrets = "secret-rotation-v2-rotate-secrets",
SecretRotationV2SendNotification = "secret-rotation-v2-send-notification",
@@ -242,10 +240,6 @@ export type TQueueJobTypes = {
| {
name: QueueJobs.PkiSyncRemoveCertificates;
payload: TQueuePkiSyncRemoveCertificatesByIdDTO;
}
| {
name: QueueJobs.PkiSyncSendActionFailedNotifications;
payload: TQueueSendPkiSyncActionFailedNotificationsDTO;
};
[QueueName.ProjectV3Migration]: {
name: QueueJobs.ProjectV3Migration;

View File

@@ -1844,52 +1844,6 @@ export const registerRoutes = async (
licenseService
});
const certificateAuthorityQueue = certificateAuthorityQueueFactory({
certificateAuthorityCrlDAL,
certificateAuthorityDAL,
certificateAuthoritySecretDAL,
certificateDAL,
projectDAL,
kmsService,
queueService,
pkiSubscriberDAL,
certificateBodyDAL,
certificateSecretDAL,
externalCertificateAuthorityDAL,
keyStore,
appConnectionDAL,
appConnectionService
});
const internalCertificateAuthorityService = internalCertificateAuthorityServiceFactory({
certificateAuthorityDAL,
certificateAuthorityCertDAL,
certificateAuthoritySecretDAL,
certificateAuthorityCrlDAL,
certificateTemplateDAL,
certificateAuthorityQueue,
certificateDAL,
certificateBodyDAL,
certificateSecretDAL,
pkiCollectionDAL,
pkiCollectionItemDAL,
projectDAL,
internalCertificateAuthorityDAL,
kmsService,
permissionService
});
const certificateEstService = certificateEstServiceFactory({
internalCertificateAuthorityService,
certificateTemplateService,
certificateTemplateDAL,
certificateAuthorityCertDAL,
certificateAuthorityDAL,
projectDAL,
kmsService,
licenseService
});
const kmipService = kmipServiceFactory({
kmipClientDAL,
permissionService,
@@ -1932,6 +1886,71 @@ export const registerRoutes = async (
gatewayV2Service
});
const pkiSyncQueue = pkiSyncQueueFactory({
queueService,
kmsService,
appConnectionDAL,
keyStore,
pkiSyncDAL,
auditLogService,
projectDAL,
licenseService,
certificateDAL,
certificateBodyDAL,
certificateSecretDAL
});
const internalCaFns = InternalCertificateAuthorityFns({
certificateAuthorityDAL,
certificateAuthorityCertDAL,
certificateAuthoritySecretDAL,
certificateAuthorityCrlDAL,
certificateDAL,
certificateBodyDAL,
certificateSecretDAL,
projectDAL,
kmsService,
pkiSyncDAL,
pkiSyncQueue
});
const certificateAuthorityQueue = certificateAuthorityQueueFactory({
certificateAuthorityCrlDAL,
certificateAuthorityDAL,
certificateAuthoritySecretDAL,
certificateDAL,
projectDAL,
kmsService,
queueService,
pkiSubscriberDAL,
certificateBodyDAL,
certificateSecretDAL,
externalCertificateAuthorityDAL,
keyStore,
appConnectionDAL,
appConnectionService,
pkiSyncDAL,
pkiSyncQueue
});
const internalCertificateAuthorityService = internalCertificateAuthorityServiceFactory({
certificateAuthorityDAL,
certificateAuthorityCertDAL,
certificateAuthoritySecretDAL,
certificateAuthorityCrlDAL,
certificateTemplateDAL,
certificateAuthorityQueue,
certificateDAL,
certificateBodyDAL,
certificateSecretDAL,
pkiCollectionDAL,
pkiCollectionItemDAL,
projectDAL,
internalCertificateAuthorityDAL,
kmsService,
permissionService
});
const certificateAuthorityService = certificateAuthorityServiceFactory({
certificateAuthorityDAL,
permissionService,
@@ -1944,19 +1963,20 @@ export const registerRoutes = async (
certificateSecretDAL,
kmsService,
pkiSubscriberDAL,
projectDAL
projectDAL,
pkiSyncDAL,
pkiSyncQueue
});
const internalCaFns = InternalCertificateAuthorityFns({
certificateAuthorityDAL,
const certificateEstService = certificateEstServiceFactory({
internalCertificateAuthorityService,
certificateTemplateService,
certificateTemplateDAL,
certificateAuthorityCertDAL,
certificateAuthoritySecretDAL,
certificateAuthorityCrlDAL,
certificateDAL,
certificateBodyDAL,
certificateSecretDAL,
certificateAuthorityDAL,
projectDAL,
kmsService
kmsService,
licenseService
});
const pkiSubscriberQueue = pkiSubscriberQueueServiceFactory({
@@ -1969,21 +1989,6 @@ export const registerRoutes = async (
internalCaFns
});
const pkiSyncQueue = pkiSyncQueueFactory({
queueService,
kmsService,
appConnectionDAL,
keyStore,
pkiSyncDAL,
auditLogService,
projectMembershipDAL,
projectDAL,
licenseService,
certificateDAL,
certificateBodyDAL,
certificateSecretDAL
});
const certificateService = certificateServiceFactory({
certificateDAL,
certificateBodyDAL,

View File

@@ -148,6 +148,15 @@ export const registerV1Routes = async (server: FastifyZodProvider) => {
await pkiRouter.register(registerPkiAlertRouter, { prefix: "/alerts" });
await pkiRouter.register(registerPkiCollectionRouter, { prefix: "/collections" });
await pkiRouter.register(registerPkiSubscriberRouter, { prefix: "/subscribers" });
await pkiRouter.register(
async (pkiSyncRouter) => {
await pkiSyncRouter.register(registerPkiSyncRouter);
for await (const [destination, router] of Object.entries(PKI_SYNC_REGISTER_ROUTER_MAP)) {
await pkiSyncRouter.register(router, { prefix: `/${destination}` });
}
},
{ prefix: "/syncs" }
);
},
{ prefix: "/pki" }
);
@@ -157,16 +166,6 @@ export const registerV1Routes = async (server: FastifyZodProvider) => {
await server.register(registerIntegrationAuthRouter, { prefix: "/integration-auth" });
await server.register(registerWebhookRouter, { prefix: "/webhooks" });
await server.register(registerIdentityRouter, { prefix: "/identities" });
await server.register(
async (pkiSyncRouter) => {
// register generic pki sync endpoints
await pkiSyncRouter.register(registerPkiSyncRouter);
for await (const [destination, router] of Object.entries(PKI_SYNC_REGISTER_ROUTER_MAP)) {
await pkiSyncRouter.register(router, { prefix: `/${destination}` });
}
},
{ prefix: "/pki-syncs" }
);
await server.register(
async (secretSharingRouter) => {

View File

@@ -1,4 +1,5 @@
import {
AZURE_KEY_VAULT_PKI_SYNC_LIST_OPTION,
AzureKeyVaultPkiSyncSchema,
CreateAzureKeyVaultPkiSyncSchema,
UpdateAzureKeyVaultPkiSyncSchema
@@ -13,5 +14,9 @@ export const registerAzureKeyVaultPkiSyncRouter = async (server: FastifyZodProvi
server,
responseSchema: AzureKeyVaultPkiSyncSchema,
createSchema: CreateAzureKeyVaultPkiSyncSchema,
updateSchema: UpdateAzureKeyVaultPkiSyncSchema
updateSchema: UpdateAzureKeyVaultPkiSyncSchema,
syncOptions: {
canImportCertificates: AZURE_KEY_VAULT_PKI_SYNC_LIST_OPTION.canImportCertificates,
canRemoveCertificates: AZURE_KEY_VAULT_PKI_SYNC_LIST_OPTION.canRemoveCertificates
}
});

View File

@@ -13,7 +13,8 @@ export const registerSyncPkiEndpoints = ({
destination,
createSchema,
updateSchema,
responseSchema
responseSchema,
syncOptions
}: {
destination: PkiSync;
server: FastifyZodProvider;
@@ -37,6 +38,10 @@ export const registerSyncPkiEndpoints = ({
subscriberId?: string;
}>;
responseSchema: z.ZodTypeAny;
syncOptions: {
canImportCertificates: boolean;
canRemoveCertificates: boolean;
};
}) => {
const destinationName = PKI_SYNC_NAME_MAP[destination];
@@ -57,7 +62,7 @@ export const registerSyncPkiEndpoints = ({
200: z.object({ pkiSyncs: responseSchema.array() })
}
},
onRequest: verifyAuth([AuthMode.JWT, AuthMode.API_KEY, AuthMode.SERVICE_TOKEN]),
onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]),
handler: async (req) => {
const {
query: { projectId }
@@ -93,19 +98,15 @@ export const registerSyncPkiEndpoints = ({
params: z.object({
pkiSyncId: z.string()
}),
querystring: z.object({
projectId: z.string().trim().min(1)
}),
response: {
200: z.object({ pkiSync: responseSchema })
200: responseSchema
}
},
onRequest: verifyAuth([AuthMode.JWT, AuthMode.API_KEY, AuthMode.SERVICE_TOKEN]),
onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]),
handler: async (req) => {
const { pkiSyncId } = req.params;
const { projectId } = req.query;
const pkiSync = await server.services.pkiSync.findPkiSyncById({ id: pkiSyncId, projectId }, req.permission);
const pkiSync = await server.services.pkiSync.findPkiSyncById({ id: pkiSyncId }, req.permission);
await server.services.auditLog.createAuditLog({
...req.auditLogInfo,
@@ -119,7 +120,7 @@ export const registerSyncPkiEndpoints = ({
}
});
return { pkiSync };
return pkiSync;
}
});
@@ -135,10 +136,10 @@ export const registerSyncPkiEndpoints = ({
description: `Create a ${destinationName} PKI Sync for the specified project.`,
body: createSchema,
response: {
200: z.object({ pkiSync: responseSchema })
200: responseSchema
}
},
onRequest: verifyAuth([AuthMode.JWT, AuthMode.API_KEY, AuthMode.SERVICE_TOKEN]),
onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]),
handler: async (req) => {
const pkiSync = await server.services.pkiSync.createPkiSync({ ...req.body, destination }, req.permission);
@@ -155,7 +156,7 @@ export const registerSyncPkiEndpoints = ({
}
});
return { pkiSync };
return pkiSync;
}
});
@@ -172,27 +173,20 @@ export const registerSyncPkiEndpoints = ({
params: z.object({
pkiSyncId: z.string()
}),
querystring: z.object({
projectId: z.string().trim().min(1)
}),
body: updateSchema,
response: {
200: z.object({ pkiSync: responseSchema })
200: responseSchema
}
},
onRequest: verifyAuth([AuthMode.JWT, AuthMode.API_KEY, AuthMode.SERVICE_TOKEN]),
onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]),
handler: async (req) => {
const { pkiSyncId } = req.params;
const { projectId } = req.query;
const pkiSync = await server.services.pkiSync.updatePkiSync(
{ ...req.body, id: pkiSyncId, projectId },
req.permission
);
const pkiSync = await server.services.pkiSync.updatePkiSync({ ...req.body, id: pkiSyncId }, req.permission);
await server.services.auditLog.createAuditLog({
...req.auditLogInfo,
projectId,
projectId: pkiSync.projectId,
event: {
type: EventType.UPDATE_PKI_SYNC,
metadata: {
@@ -202,7 +196,7 @@ export const registerSyncPkiEndpoints = ({
}
});
return { pkiSync };
return pkiSync;
}
});
@@ -219,23 +213,19 @@ export const registerSyncPkiEndpoints = ({
params: z.object({
pkiSyncId: z.string()
}),
querystring: z.object({
projectId: z.string().trim().min(1)
}),
response: {
200: z.object({ pkiSync: responseSchema })
200: responseSchema
}
},
onRequest: verifyAuth([AuthMode.JWT, AuthMode.API_KEY, AuthMode.SERVICE_TOKEN]),
onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]),
handler: async (req) => {
const { pkiSyncId } = req.params;
const { projectId } = req.query;
const pkiSync = await server.services.pkiSync.deletePkiSync({ id: pkiSyncId, projectId }, req.permission);
const pkiSync = await server.services.pkiSync.deletePkiSync({ id: pkiSyncId }, req.permission);
await server.services.auditLog.createAuditLog({
...req.auditLogInfo,
projectId,
projectId: pkiSync.projectId,
event: {
type: EventType.DELETE_PKI_SYNC,
metadata: {
@@ -246,7 +236,7 @@ export const registerSyncPkiEndpoints = ({
}
});
return { pkiSync };
return pkiSync;
}
});
@@ -263,22 +253,17 @@ export const registerSyncPkiEndpoints = ({
params: z.object({
pkiSyncId: z.string()
}),
querystring: z.object({
projectId: z.string().trim().min(1)
}),
response: {
200: z.object({ message: z.string() })
}
},
onRequest: verifyAuth([AuthMode.JWT, AuthMode.API_KEY, AuthMode.SERVICE_TOKEN]),
onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]),
handler: async (req) => {
const { pkiSyncId } = req.params;
const { projectId } = req.query;
const result = await server.services.pkiSync.triggerPkiSyncSyncCertificatesById(
{
id: pkiSyncId,
projectId
id: pkiSyncId
},
req.permission
);
@@ -287,42 +272,40 @@ export const registerSyncPkiEndpoints = ({
}
});
server.route({
method: "POST",
url: "/:pkiSyncId/import",
config: {
rateLimit: writeLimit
},
schema: {
hide: false,
tags: [ApiDocsTags.PkiSyncs],
description: `Import certificates from the specified ${destinationName} PKI Sync destination.`,
params: z.object({
pkiSyncId: z.string()
}),
querystring: z.object({
projectId: z.string().trim().min(1)
}),
response: {
200: z.object({ message: z.string() })
// Only register import route if the destination supports it
if (syncOptions.canImportCertificates) {
server.route({
method: "POST",
url: "/:pkiSyncId/import",
config: {
rateLimit: writeLimit
},
schema: {
hide: false,
tags: [ApiDocsTags.PkiSyncs],
description: `Import certificates from the specified ${destinationName} PKI Sync destination.`,
params: z.object({
pkiSyncId: z.string()
}),
response: {
200: z.object({ message: z.string() })
}
},
onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]),
handler: async (req) => {
const { pkiSyncId } = req.params;
const result = await server.services.pkiSync.triggerPkiSyncImportCertificatesById(
{
id: pkiSyncId
},
req.permission
);
return result;
}
},
onRequest: verifyAuth([AuthMode.JWT, AuthMode.API_KEY, AuthMode.SERVICE_TOKEN]),
handler: async (req) => {
const { pkiSyncId } = req.params;
const { projectId } = req.query;
const result = await server.services.pkiSync.triggerPkiSyncImportCertificatesById(
{
id: pkiSyncId,
projectId
},
req.permission
);
return result;
}
});
});
}
server.route({
method: "POST",
@@ -337,22 +320,17 @@ export const registerSyncPkiEndpoints = ({
params: z.object({
pkiSyncId: z.string()
}),
querystring: z.object({
projectId: z.string().trim().min(1)
}),
response: {
200: z.object({ message: z.string() })
}
},
onRequest: verifyAuth([AuthMode.JWT, AuthMode.API_KEY, AuthMode.SERVICE_TOKEN]),
onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]),
handler: async (req) => {
const { pkiSyncId } = req.params;
const { projectId } = req.query;
const result = await server.services.pkiSync.triggerPkiSyncRemoveCertificatesById(
{
id: pkiSyncId,
projectId
id: pkiSyncId
},
req.permission
);

View File

@@ -61,7 +61,12 @@ const PkiSyncOptionsSchema = z.object({
connection: z.nativeEnum(AppConnection),
destination: z.nativeEnum(PkiSync),
canImportCertificates: z.boolean(),
canRemoveCertificates: z.boolean()
canRemoveCertificates: z.boolean(),
defaultCertificateNameSchema: z.string().optional(),
forbiddenCharacters: z.string().optional(),
allowedCharacterPattern: z.string().optional(),
maxCertificateNameLength: z.number().optional(),
minCertificateNameLength: z.number().optional()
});
export const registerPkiSyncRouter = async (server: FastifyZodProvider) => {
@@ -81,7 +86,7 @@ export const registerPkiSyncRouter = async (server: FastifyZodProvider) => {
})
}
},
onRequest: verifyAuth([AuthMode.JWT, AuthMode.API_KEY, AuthMode.SERVICE_TOKEN]),
onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]),
handler: () => {
const pkiSyncOptions = server.services.pkiSync.getPkiSyncOptions();
return { pkiSyncOptions };
@@ -105,7 +110,7 @@ export const registerPkiSyncRouter = async (server: FastifyZodProvider) => {
200: z.object({ pkiSyncs: PkiSyncSchema.array() })
}
},
onRequest: verifyAuth([AuthMode.JWT, AuthMode.API_KEY, AuthMode.SERVICE_TOKEN]),
onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]),
handler: async (req) => {
const {
query: { projectId },
@@ -142,19 +147,15 @@ export const registerPkiSyncRouter = async (server: FastifyZodProvider) => {
params: z.object({
pkiSyncId: z.string()
}),
querystring: z.object({
projectId: z.string().trim().min(1)
}),
response: {
200: z.object({ pkiSync: PkiSyncSchema })
200: PkiSyncSchema
}
},
onRequest: verifyAuth([AuthMode.JWT, AuthMode.API_KEY, AuthMode.SERVICE_TOKEN]),
onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]),
handler: async (req) => {
const { pkiSyncId } = req.params;
const { projectId } = req.query;
const pkiSync = await server.services.pkiSync.findPkiSyncById({ id: pkiSyncId, projectId }, req.permission);
const pkiSync = await server.services.pkiSync.findPkiSyncById({ id: pkiSyncId }, req.permission);
await server.services.auditLog.createAuditLog({
...req.auditLogInfo,
@@ -168,7 +169,7 @@ export const registerPkiSyncRouter = async (server: FastifyZodProvider) => {
}
});
return { pkiSync };
return pkiSync;
}
});
};

View File

@@ -23,6 +23,7 @@ import {
} from "@app/services/certificate/certificate-types";
import { TKmsServiceFactory } from "@app/services/kms/kms-service";
import { TPkiSubscriberDALFactory } from "@app/services/pki-subscriber/pki-subscriber-dal";
import { triggerAutoSyncForSubscriber } from "@app/services/pki-sync/pki-sync-fns";
import { TProjectDALFactory } from "@app/services/project/project-dal";
import { getProjectKmsCertificateKeyId } from "@app/services/project/project-fns";
@@ -56,6 +57,12 @@ type TAcmeCertificateAuthorityFnsDeps = {
"encryptWithKmsKey" | "generateKmsKey" | "createCipherPairWithDataKey" | "decryptWithKmsKey"
>;
pkiSubscriberDAL: Pick<TPkiSubscriberDALFactory, "findById">;
pkiSyncDAL: {
find: (filter: { subscriberId: string; isAutoSyncEnabled: boolean }) => Promise<Array<{ id: string }>>;
};
pkiSyncQueue: {
queuePkiSyncSyncCertificatesById: (params: { syncId: string }) => Promise<void>;
};
projectDAL: Pick<TProjectDALFactory, "findById" | "findOne" | "updateById" | "transaction">;
};
@@ -109,7 +116,9 @@ export const AcmeCertificateAuthorityFns = ({
certificateSecretDAL,
kmsService,
projectDAL,
pkiSubscriberDAL
pkiSubscriberDAL,
pkiSyncDAL,
pkiSyncQueue
}: TAcmeCertificateAuthorityFnsDeps) => {
const createCertificateAuthority = async ({
name,
@@ -524,6 +533,8 @@ export const AcmeCertificateAuthorityFns = ({
tx
);
});
await triggerAutoSyncForSubscriber(subscriber.id, { pkiSyncDAL, pkiSyncQueue });
};
return {

View File

@@ -26,6 +26,7 @@ import {
import { TKmsServiceFactory } from "@app/services/kms/kms-service";
import { TPkiSubscriberDALFactory } from "@app/services/pki-subscriber/pki-subscriber-dal";
import { TPkiSubscriberProperties } from "@app/services/pki-subscriber/pki-subscriber-types";
import { triggerAutoSyncForSubscriber } from "@app/services/pki-sync/pki-sync-fns";
import { TProjectDALFactory } from "@app/services/project/project-dal";
import { getProjectKmsCertificateKeyId } from "@app/services/project/project-fns";
@@ -55,6 +56,12 @@ type TAzureAdCsCertificateAuthorityFnsDeps = {
"encryptWithKmsKey" | "generateKmsKey" | "createCipherPairWithDataKey" | "decryptWithKmsKey"
>;
pkiSubscriberDAL: Pick<TPkiSubscriberDALFactory, "findById">;
pkiSyncDAL: {
find: (filter: { subscriberId: string; isAutoSyncEnabled: boolean }) => Promise<Array<{ id: string }>>;
};
pkiSyncQueue: {
queuePkiSyncSyncCertificatesById: (params: { syncId: string }) => Promise<void>;
};
projectDAL: Pick<TProjectDALFactory, "findById" | "findOne" | "updateById" | "transaction">;
};
@@ -584,7 +591,9 @@ export const AzureAdCsCertificateAuthorityFns = ({
certificateSecretDAL,
kmsService,
projectDAL,
pkiSubscriberDAL
pkiSubscriberDAL,
pkiSyncDAL,
pkiSyncQueue
}: TAzureAdCsCertificateAuthorityFnsDeps) => {
const createCertificateAuthority = async ({
name,
@@ -1024,6 +1033,8 @@ export const AzureAdCsCertificateAuthorityFns = ({
);
});
await triggerAutoSyncForSubscriber(subscriber.id, { pkiSyncDAL, pkiSyncQueue });
return {
certificate: certificatePem,
certificateChain: certificateChainPem,

View File

@@ -50,6 +50,12 @@ type TCertificateAuthorityQueueFactoryDep = {
certificateSecretDAL: Pick<TCertificateSecretDALFactory, "create">;
queueService: TQueueServiceFactory;
pkiSubscriberDAL: Pick<TPkiSubscriberDALFactory, "findById" | "updateById">;
pkiSyncDAL: {
find: (filter: { subscriberId: string; isAutoSyncEnabled: boolean }) => Promise<Array<{ id: string }>>;
};
pkiSyncQueue: {
queuePkiSyncSyncCertificatesById: (params: { syncId: string }) => Promise<void>;
};
};
export type TCertificateAuthorityQueueFactory = ReturnType<typeof certificateAuthorityQueueFactory>;
@@ -68,7 +74,9 @@ export const certificateAuthorityQueueFactory = ({
externalCertificateAuthorityDAL,
certificateBodyDAL,
certificateSecretDAL,
pkiSubscriberDAL
pkiSubscriberDAL,
pkiSyncDAL,
pkiSyncQueue
}: TCertificateAuthorityQueueFactoryDep) => {
const acmeFns = AcmeCertificateAuthorityFns({
appConnectionDAL,
@@ -80,7 +88,9 @@ export const certificateAuthorityQueueFactory = ({
certificateSecretDAL,
kmsService,
pkiSubscriberDAL,
projectDAL
projectDAL,
pkiSyncDAL,
pkiSyncQueue
});
const azureAdCsFns = AzureAdCsCertificateAuthorityFns({
@@ -93,7 +103,9 @@ export const certificateAuthorityQueueFactory = ({
certificateSecretDAL,
kmsService,
pkiSubscriberDAL,
projectDAL
projectDAL,
pkiSyncDAL,
pkiSyncQueue
});
// TODO 1: auto-periodic rotation

View File

@@ -68,6 +68,12 @@ type TCertificateAuthorityServiceFactoryDep = {
"encryptWithKmsKey" | "generateKmsKey" | "createCipherPairWithDataKey" | "decryptWithKmsKey"
>;
pkiSubscriberDAL: Pick<TPkiSubscriberDALFactory, "findById">;
pkiSyncDAL: {
find: (filter: { subscriberId: string; isAutoSyncEnabled: boolean }) => Promise<Array<{ id: string }>>;
};
pkiSyncQueue: {
queuePkiSyncSyncCertificatesById: (params: { syncId: string }) => Promise<void>;
};
};
export type TCertificateAuthorityServiceFactory = ReturnType<typeof certificateAuthorityServiceFactory>;
@@ -84,7 +90,9 @@ export const certificateAuthorityServiceFactory = ({
certificateBodyDAL,
certificateSecretDAL,
kmsService,
pkiSubscriberDAL
pkiSubscriberDAL,
pkiSyncDAL,
pkiSyncQueue
}: TCertificateAuthorityServiceFactoryDep) => {
const acmeFns = AcmeCertificateAuthorityFns({
appConnectionDAL,
@@ -96,7 +104,9 @@ export const certificateAuthorityServiceFactory = ({
certificateSecretDAL,
kmsService,
pkiSubscriberDAL,
projectDAL
projectDAL,
pkiSyncDAL,
pkiSyncQueue
});
const azureAdCsFns = AzureAdCsCertificateAuthorityFns({
@@ -109,7 +119,9 @@ export const certificateAuthorityServiceFactory = ({
certificateSecretDAL,
kmsService,
pkiSubscriberDAL,
projectDAL
projectDAL,
pkiSyncDAL,
pkiSyncQueue
});
const createCertificateAuthority = async (

View File

@@ -19,6 +19,7 @@ import {
TAltNameMapping
} from "@app/services/certificate/certificate-types";
import { TKmsServiceFactory } from "@app/services/kms/kms-service";
import { triggerAutoSyncForSubscriber } from "@app/services/pki-sync/pki-sync-fns";
import { TProjectDALFactory } from "@app/services/project/project-dal";
import { getProjectKmsCertificateKeyId } from "@app/services/project/project-fns";
@@ -51,6 +52,12 @@ type TInternalCertificateAuthorityFnsDeps = {
certificateDAL: Pick<TCertificateDALFactory, "create" | "transaction">;
certificateBodyDAL: Pick<TCertificateBodyDALFactory, "create">;
certificateSecretDAL: Pick<TCertificateSecretDALFactory, "create">;
pkiSyncDAL: {
find: (filter: { subscriberId: string; isAutoSyncEnabled: boolean }) => Promise<Array<{ id: string }>>;
};
pkiSyncQueue: {
queuePkiSyncSyncCertificatesById: (params: { syncId: string }) => Promise<void>;
};
};
export const InternalCertificateAuthorityFns = ({
@@ -62,7 +69,9 @@ export const InternalCertificateAuthorityFns = ({
certificateAuthorityCrlDAL,
certificateDAL,
certificateBodyDAL,
certificateSecretDAL
certificateSecretDAL,
pkiSyncDAL,
pkiSyncQueue
}: TInternalCertificateAuthorityFnsDeps) => {
const issueCertificate = async (
subscriber: TPkiSubscribers,
@@ -251,6 +260,8 @@ export const InternalCertificateAuthorityFns = ({
);
});
await triggerAutoSyncForSubscriber(subscriber.id, { pkiSyncDAL, pkiSyncQueue });
return {
certificate: leafCert.toString("pem"),
certificateChain: certificateChainPem,

View File

@@ -11,7 +11,6 @@ import {
} from "@app/ee/services/permission/project-permission";
import { crypto } from "@app/lib/crypto/cryptography";
import { BadRequestError, NotFoundError } from "@app/lib/errors";
import { logger } from "@app/lib/logger";
import { TCertificateBodyDALFactory } from "@app/services/certificate/certificate-body-dal";
import { TCertificateDALFactory } from "@app/services/certificate/certificate-dal";
import { TCertificateAuthorityCertDALFactory } from "@app/services/certificate-authority/certificate-authority-cert-dal";
@@ -23,6 +22,7 @@ import { TKmsServiceFactory } from "@app/services/kms/kms-service";
import { TPkiCollectionDALFactory } from "@app/services/pki-collection/pki-collection-dal";
import { TPkiCollectionItemDALFactory } from "@app/services/pki-collection/pki-collection-item-dal";
import { TPkiSyncDALFactory } from "@app/services/pki-sync/pki-sync-dal";
import { triggerAutoSyncForSubscriber } from "@app/services/pki-sync/pki-sync-fns";
import { TPkiSyncQueueFactory } from "@app/services/pki-sync/pki-sync-queue";
import { TProjectDALFactory } from "@app/services/project/project-dal";
import { getProjectKmsCertificateKeyId } from "@app/services/project/project-fns";
@@ -79,28 +79,6 @@ export const certificateServiceFactory = ({
pkiSyncDAL,
pkiSyncQueue
}: TCertificateServiceFactoryDep) => {
/**
* Trigger auto sync for PKI syncs connected to a PKI subscriber when certificates are issued/revoked/deleted
*/
const triggerAutoSyncForSubscriber = async (subscriberId: string) => {
try {
// Find all PKI syncs that are connected to this subscriber and have auto sync enabled
const pkiSyncs = await pkiSyncDAL.find({
subscriberId,
isAutoSyncEnabled: true
});
// Queue sync jobs for each auto sync enabled PKI sync
for (const pkiSync of pkiSyncs) {
await pkiSyncQueue.queuePkiSyncSyncCertificatesById({ syncId: pkiSync.id });
}
} catch (error) {
// Don't throw error to avoid breaking the main certificate operation
// Just log the auto sync failure
logger.error(error, `Failed to trigger auto sync for subscriber ${subscriberId}:`);
}
};
/**
* Return details for certificate with serial number [serialNumber]
*/
@@ -190,7 +168,7 @@ export const certificateServiceFactory = ({
// Trigger auto sync for PKI syncs connected to this certificate's subscriber
if (cert.pkiSubscriberId) {
await triggerAutoSyncForSubscriber(cert.pkiSubscriberId);
await triggerAutoSyncForSubscriber(cert.pkiSubscriberId, { pkiSyncDAL, pkiSyncQueue });
}
return {
@@ -259,7 +237,7 @@ export const certificateServiceFactory = ({
// Trigger auto sync for PKI syncs connected to this certificate's subscriber
if (cert.pkiSubscriberId) {
await triggerAutoSyncForSubscriber(cert.pkiSubscriberId);
await triggerAutoSyncForSubscriber(cert.pkiSubscriberId, { pkiSyncDAL, pkiSyncQueue });
}
// Note: External CA revocation handling would go here for supported CA types

View File

@@ -0,0 +1,129 @@
/* eslint-disable no-await-in-loop */
import { AxiosError } from "axios";
import { logger } from "@app/lib/logger";
export type RateLimitConfig = {
MAX_CONCURRENT_REQUESTS: number;
BASE_DELAY: number;
MAX_DELAY: number;
MAX_RETRIES: number;
RATE_LIMIT_STATUS_CODES: number[];
};
export type RateLimitContext = {
operation: string;
identifier?: string;
syncId: string;
};
export type ConcurrencyContext = {
operation: string;
syncId: string;
};
export const sleep = (ms: number): Promise<void> =>
new Promise((resolve) => {
setTimeout(resolve, ms);
});
export const createRateLimitErrorChecker =
(config: RateLimitConfig) =>
(error: unknown): boolean => {
if (error instanceof AxiosError) {
return (
config.RATE_LIMIT_STATUS_CODES.includes(error.response?.status || 0) ||
error.message.toLowerCase().includes("rate limit") ||
error.message.toLowerCase().includes("throttl")
);
}
return false;
};
export const createRateLimitRetry =
(config: RateLimitConfig, isRateLimitError: (error: unknown) => boolean) =>
async <T>(fn: () => Promise<T>, context: RateLimitContext, retryCount = 0): Promise<T> => {
try {
return await fn();
} catch (error) {
if (isRateLimitError(error) && retryCount < config.MAX_RETRIES) {
const delay = Math.min(config.BASE_DELAY * 2 ** retryCount, config.MAX_DELAY);
logger.warn(
{
syncId: context.syncId,
operation: context.operation,
identifier: context.identifier,
retryCount: retryCount + 1,
delayMs: delay,
error: error instanceof AxiosError ? error.message : String(error)
},
"Rate limit hit, retrying with exponential backoff"
);
await sleep(delay);
return createRateLimitRetry(config, isRateLimitError)(fn, context, retryCount + 1);
}
throw error;
}
};
export const createConcurrencyLimitExecutor =
(
config: RateLimitConfig,
withRateLimitRetry: <T>(fn: () => Promise<T>, context: RateLimitContext, retryCount?: number) => Promise<T>
) =>
async <T, R>(
items: T[],
executor: (item: T) => Promise<R>,
context: ConcurrencyContext,
concurrencyLimit = config.MAX_CONCURRENT_REQUESTS
): Promise<PromiseSettledResult<R>[]> => {
const results: PromiseSettledResult<R>[] = [];
for (let i = 0; i < items.length; i += concurrencyLimit) {
const batch = items.slice(i, i + concurrencyLimit);
logger.debug(
{
syncId: context.syncId,
operation: context.operation,
batchStart: i + 1,
batchEnd: Math.min(i + concurrencyLimit, items.length),
totalItems: items.length
},
"Processing batch with rate limit protection"
);
const batchPromises = batch.map((item, batchIndex) =>
withRateLimitRetry(() => executor(item), {
operation: context.operation,
identifier: `batch-${i + batchIndex + 1}`,
syncId: context.syncId
})
);
const batchResults = await Promise.allSettled(batchPromises);
results.push(...batchResults);
if (i + concurrencyLimit < items.length) {
await sleep(100);
}
}
return results;
};
export const createConnectionQueue = (config: RateLimitConfig) => {
const isRateLimitError = createRateLimitErrorChecker(config);
const withRateLimitRetry = createRateLimitRetry(config, isRateLimitError);
const executeWithConcurrencyLimit = createConcurrencyLimitExecutor(config, withRateLimitRetry);
return {
sleep,
isRateLimitError,
withRateLimitRetry,
executeWithConcurrencyLimit
};
};

View File

@@ -0,0 +1 @@
export * from "./connection-queue-fns";

View File

@@ -13,7 +13,6 @@ import {
} from "@app/ee/services/permission/project-permission";
import { getConfig } from "@app/lib/config/env";
import { BadRequestError, NotFoundError } from "@app/lib/errors";
import { logger } from "@app/lib/logger";
import { ms } from "@app/lib/ms";
import { TCertificateBodyDALFactory } from "@app/services/certificate/certificate-body-dal";
import { TCertificateDALFactory } from "@app/services/certificate/certificate-dal";
@@ -39,6 +38,7 @@ import { TCertificateAuthoritySecretDALFactory } from "@app/services/certificate
import { TKmsServiceFactory } from "@app/services/kms/kms-service";
import { TPkiSubscriberDALFactory } from "@app/services/pki-subscriber/pki-subscriber-dal";
import { TPkiSyncDALFactory } from "@app/services/pki-sync/pki-sync-dal";
import { triggerAutoSyncForSubscriber } from "@app/services/pki-sync/pki-sync-fns";
import { TPkiSyncQueueFactory } from "@app/services/pki-sync/pki-sync-queue";
import { TProjectDALFactory } from "@app/services/project/project-dal";
import { getProjectKmsCertificateKeyId } from "@app/services/project/project-fns";
@@ -106,28 +106,6 @@ export const pkiSubscriberServiceFactory = ({
pkiSyncDAL,
pkiSyncQueue
}: TPkiSubscriberServiceFactoryDep) => {
/**
* Trigger auto sync for PKI syncs connected to a PKI subscriber when certificates are issued
*/
const triggerAutoSyncForSubscriber = async (subscriberId: string) => {
try {
// Find all PKI syncs that are connected to this subscriber and have auto sync enabled
const pkiSyncs = await pkiSyncDAL.find({
subscriberId,
isAutoSyncEnabled: true
});
// Queue sync jobs for each auto sync enabled PKI sync
for (const pkiSync of pkiSyncs) {
await pkiSyncQueue.queuePkiSyncSyncCertificatesById({ syncId: pkiSync.id });
}
} catch (error) {
// Don't throw error to avoid breaking the main certificate operation
// Just log the auto sync failure
logger.error(error, `Failed to trigger auto sync for subscriber ${subscriberId}:`);
}
};
const createSubscriber = async ({
name,
commonName,
@@ -446,7 +424,7 @@ export const pkiSubscriberServiceFactory = ({
const result = await internalCaFns.issueCertificate(subscriber, ca);
// Trigger auto sync for PKI syncs connected to this subscriber after certificate issuance
await triggerAutoSyncForSubscriber(subscriber.id);
await triggerAutoSyncForSubscriber(subscriber.id, { pkiSyncDAL, pkiSyncQueue });
return result;
}
@@ -707,7 +685,7 @@ export const pkiSubscriberServiceFactory = ({
});
// Trigger auto sync for PKI syncs connected to this subscriber after certificate signing
await triggerAutoSyncForSubscriber(subscriber.id);
await triggerAutoSyncForSubscriber(subscriber.id, { pkiSyncDAL, pkiSyncQueue });
return {
certificate: leafCert.toString("pem"),

View File

@@ -1,24 +1,61 @@
/* eslint-disable no-await-in-loop */
import { AxiosError } from "axios";
import * as crypto from "crypto";
import { request } from "@app/lib/config/request";
import { logger } from "@app/lib/logger";
import { TAppConnectionDALFactory } from "@app/services/app-connection/app-connection-dal";
import { AppConnection } from "@app/services/app-connection/app-connection-enums";
import { getAzureConnectionAccessToken } from "@app/services/app-connection/azure-key-vault";
import { createConnectionQueue, RateLimitConfig } from "@app/services/connection-queue";
import { TKmsServiceFactory } from "@app/services/kms/kms-service";
import { matchesCertificateNameSchema } from "@app/services/pki-sync/pki-sync-fns";
import { TCertificateMap } from "@app/services/pki-sync/pki-sync-types";
import { PkiSync } from "../pki-sync-enums";
import { PkiSyncError } from "../pki-sync-errors";
import { GetAzureKeyVaultCertificate, TAzureKeyVaultPkiSyncWithCredentials } from "./azure-key-vault-pki-sync-types";
import { TPkiSyncWithCredentials } from "../pki-sync-types";
import { GetAzureKeyVaultCertificate, TAzureKeyVaultPkiSyncConfig } from "./azure-key-vault-pki-sync-types";
const AZURE_RATE_LIMIT_CONFIG: RateLimitConfig = {
MAX_CONCURRENT_REQUESTS: 10,
BASE_DELAY: 1000,
MAX_DELAY: 30000,
MAX_RETRIES: 3,
RATE_LIMIT_STATUS_CODES: [429, 503]
};
const azureConnectionQueue = createConnectionQueue(AZURE_RATE_LIMIT_CONFIG);
const { withRateLimitRetry, executeWithConcurrencyLimit } = azureConnectionQueue;
const extractCertificateNameFromId = (certificateId: string): string => {
return certificateId.substring(certificateId.lastIndexOf("/") + 1);
};
const isInfisicalManagedCertificate = (certificateName: string, pkiSync: TPkiSyncWithCredentials): boolean => {
const syncOptions = pkiSync.syncOptions as { certificateNameSchema?: string } | undefined;
const certificateNameSchema = syncOptions?.certificateNameSchema;
if (certificateNameSchema) {
const environment = "global";
return matchesCertificateNameSchema(certificateName, environment, certificateNameSchema);
}
return certificateName.startsWith("Infisical-PKI-Sync-");
};
export const AZURE_KEY_VAULT_PKI_SYNC_LIST_OPTION = {
name: "Azure Key Vault" as const,
connection: AppConnection.AzureKeyVault,
destination: PkiSync.AzureKeyVault,
canImportCertificates: false,
canRemoveCertificates: true
canRemoveCertificates: true,
defaultCertificateNameSchema: "Infisical-PKI-Sync-{{certificateId}}",
forbiddenCharacters: "!@#$%^&*()+=[]{}|\\:;\"'<>,.?/~` _",
allowedCharacterPattern: "^[a-zA-Z0-9-]{1,127}$",
maxCertificateNameLength: 127,
minCertificateNameLength: 1
};
type TAzureKeyVaultPkiSyncFactoryDeps = {
@@ -26,19 +63,164 @@ type TAzureKeyVaultPkiSyncFactoryDeps = {
kmsService: Pick<TKmsServiceFactory, "createCipherPairWithDataKey">;
};
const parseCertificateX509Props = (certPem: string) => {
try {
const cert = new crypto.X509Certificate(certPem);
const { subject } = cert;
const sans = {
dns_names: [] as string[],
emails: [] as string[],
upns: [] as string[]
};
if (cert.subjectAltName) {
const sanEntries = cert.subjectAltName.split(", ");
for (const entry of sanEntries) {
if (entry.startsWith("DNS:")) {
sans.dns_names.push(entry.substring(4));
} else if (entry.startsWith("email:")) {
sans.emails.push(entry.substring(6));
} else if (entry.startsWith("othername:UPN:")) {
sans.upns.push(entry.substring(14));
}
}
}
return {
subject,
sans
};
} catch (error) {
logger.warn(
{ error: error instanceof Error ? error.message : String(error) },
"Failed to parse certificate X.509 properties, using empty values"
);
return {
subject: "",
sans: {
dns_names: [],
emails: [],
upns: []
}
};
}
};
const parseCertificateKeyProps = (certPem: string) => {
try {
const publicKeyObject = crypto.createPublicKey(certPem);
const keyDetails = publicKeyObject.asymmetricKeyDetails;
if (!keyDetails) {
if (publicKeyObject.asymmetricKeyType === "rsa") {
const pubKeyStr = publicKeyObject.export({ type: "spki", format: "der" }).toString("hex");
const estimatedBits = pubKeyStr.length * 4;
let keySize = 2048;
if (estimatedBits >= 4000) {
keySize = 4096;
} else if (estimatedBits >= 3000) {
keySize = 3072;
} else if (estimatedBits >= 2000) {
keySize = 2048;
} else if (estimatedBits >= 1000) {
keySize = 1024;
}
return {
kty: "RSA",
key_size: keySize
};
}
if (publicKeyObject.asymmetricKeyType === "ec") {
return {
kty: "EC",
curve: "P-256"
};
}
return {
kty: "RSA",
key_size: 2048
};
}
if (publicKeyObject.asymmetricKeyType === "rsa") {
const modulusLength = keyDetails.modulusLength || 2048;
return {
kty: "RSA",
key_size: modulusLength
};
}
if (publicKeyObject.asymmetricKeyType === "ec") {
const { namedCurve } = keyDetails;
let curveName = "P-256";
switch (namedCurve) {
case "prime256v1":
case "secp256r1":
curveName = "P-256";
break;
case "secp384r1":
curveName = "P-384";
break;
case "secp521r1":
curveName = "P-521";
break;
default:
curveName = "P-256";
}
return {
kty: "EC",
curve: curveName
};
}
const keyType = publicKeyObject.asymmetricKeyType;
if (keyType && !["rsa", "ec"].includes(keyType)) {
throw new Error(`Unsupported certificate key type: ${keyType}. Azure Key Vault only supports RSA and EC keys.`);
}
logger.warn({ keyType }, "Unable to determine certificate key type, defaulting to RSA 2048");
return {
kty: "RSA",
key_size: 2048
};
} catch (error) {
logger.warn(
{ error: error instanceof Error ? error.message : String(error) },
"Failed to parse certificate key properties, defaulting to RSA 2048"
);
return {
kty: "RSA",
key_size: 2048
};
}
};
export const azureKeyVaultPkiSyncFactory = ({ kmsService, appConnectionDAL }: TAzureKeyVaultPkiSyncFactoryDeps) => {
const $getAzureKeyVaultCertificates = async (accessToken: string, vaultBaseUrl: string) => {
const $getAzureKeyVaultCertificates = async (accessToken: string, vaultBaseUrl: string, syncId = "unknown") => {
const paginateAzureKeyVaultCertificates = async () => {
let result: GetAzureKeyVaultCertificate[] = [];
let currentUrl = `${vaultBaseUrl}/certificates?api-version=7.4`;
while (currentUrl) {
const res = await request.get<{ value: GetAzureKeyVaultCertificate[]; nextLink: string }>(currentUrl, {
headers: {
Authorization: `Bearer ${accessToken}`
}
});
const urlToFetch = currentUrl; // Capture current URL to avoid loop function issue
const res = await withRateLimitRetry(
() =>
request.get<{ value: GetAzureKeyVaultCertificate[]; nextLink: string }>(urlToFetch, {
headers: {
Authorization: `Bearer ${accessToken}`
}
}),
{ operation: "list-certificates", syncId }
);
result = result.concat(res.data.value);
currentUrl = res.data.nextLink;
@@ -54,48 +236,82 @@ export const azureKeyVaultPkiSyncFactory = ({ kmsService, appConnectionDAL }: TA
// disabled certificates to skip sending updates to
const disabledAzureKeyVaultCertificateKeys = getAzureKeyVaultCertificates
.filter(({ attributes }) => !attributes.enabled)
.map((getAzureKeyVaultCertificate) => {
return getAzureKeyVaultCertificate.id.substring(getAzureKeyVaultCertificate.id.lastIndexOf("/") + 1);
});
.map((certificate) => extractCertificateNameFromId(certificate.id));
let lastSlashIndex: number;
const res = (
await Promise.all(
enabledAzureKeyVaultCertificates.map(async (getAzureKeyVaultCertificate) => {
if (!lastSlashIndex) {
lastSlashIndex = getAzureKeyVaultCertificate.id.lastIndexOf("/");
}
const azureKeyVaultCertificate = await request.get<GetAzureKeyVaultCertificate>(
`${getAzureKeyVaultCertificate.id}?api-version=7.4`,
{
headers: {
Authorization: `Bearer ${accessToken}`
}
}
);
let certPem = "";
if (azureKeyVaultCertificate.data.cer) {
try {
// Azure Key Vault stores certificate in base64 DER format
// We need to convert it to PEM format with proper headers
const base64Cert = azureKeyVaultCertificate.data.cer;
certPem = `-----BEGIN CERTIFICATE-----\n${base64Cert.match(/.{1,64}/g)?.join("\n")}\n-----END CERTIFICATE-----`;
} catch (error) {
certPem = azureKeyVaultCertificate.data.cer;
// Use rate-limited concurrent execution for fetching certificate details
const certificateResults = await executeWithConcurrencyLimit(
enabledAzureKeyVaultCertificates,
async (getAzureKeyVaultCertificate) => {
const azureKeyVaultCertificate = await request.get<GetAzureKeyVaultCertificate>(
`${getAzureKeyVaultCertificate.id}?api-version=7.4`,
{
headers: {
Authorization: `Bearer ${accessToken}`
}
}
);
return {
...azureKeyVaultCertificate.data,
key: getAzureKeyVaultCertificate.id.substring(lastSlashIndex + 1),
cert: certPem,
privateKey: "" // Private keys cannot be extracted from Azure Key Vault for security reasons
};
})
let certPem = "";
if (azureKeyVaultCertificate.data.cer) {
try {
// Azure Key Vault stores certificate in base64 DER format
// We need to convert it to PEM format with proper headers
const base64Cert = azureKeyVaultCertificate.data.cer;
const base64Lines = base64Cert.match(/.{1,64}/g);
if (!base64Lines) {
throw new Error("Failed to format base64 certificate data");
}
certPem = `-----BEGIN CERTIFICATE-----\n${base64Lines.join("\n")}\n-----END CERTIFICATE-----`;
} catch (error) {
logger.warn(
{
error: error instanceof Error ? error.message : String(error),
certificateId: getAzureKeyVaultCertificate.id
},
"Failed to convert Azure Key Vault certificate to PEM format, skipping certificate"
);
certPem = ""; // Skip this certificate if we can't convert it properly
}
}
return {
...azureKeyVaultCertificate.data,
key: extractCertificateNameFromId(getAzureKeyVaultCertificate.id),
cert: certPem,
privateKey: "" // Private keys cannot be extracted from Azure Key Vault for security reasons
};
},
{ operation: "fetch-certificate-details", syncId }
);
const successfulCertificates = certificateResults
.filter(
(
result
): result is PromiseFulfilledResult<
GetAzureKeyVaultCertificate & {
key: string;
cert: string;
privateKey: string;
}
> => result.status === "fulfilled"
)
).reduce(
.map((result) => result.value);
// Log any failures
const failedFetches = certificateResults.filter((result) => result.status === "rejected");
if (failedFetches.length > 0) {
logger.warn(
{
syncId,
failedCount: failedFetches.length,
totalCount: enabledAzureKeyVaultCertificates.length
},
"Some certificate details could not be fetched from Azure Key Vault"
);
}
const res: Record<string, { cert: string; privateKey: string }> = successfulCertificates.reduce(
(obj, certificate) => ({
...obj,
[certificate.key]: {
@@ -112,30 +328,16 @@ export const azureKeyVaultPkiSyncFactory = ({ kmsService, appConnectionDAL }: TA
};
};
const syncCertificates = async (pkiSync: TAzureKeyVaultPkiSyncWithCredentials, certificateMap: TCertificateMap) => {
logger.info(
{
syncId: pkiSync.id,
vaultUrl: pkiSync.destinationConfig.vaultBaseUrl,
certificateCount: Object.keys(certificateMap).length
},
"Starting Azure Key Vault certificate sync"
);
const syncCertificates = async (pkiSync: TPkiSyncWithCredentials, certificateMap: TCertificateMap) => {
const { accessToken } = await getAzureConnectionAccessToken(pkiSync.connection.id, appConnectionDAL, kmsService);
// Cast destination config to Azure Key Vault config
const destinationConfig = pkiSync.destinationConfig as TAzureKeyVaultPkiSyncConfig;
const { vaultCertificates, disabledAzureKeyVaultCertificateKeys } = await $getAzureKeyVaultCertificates(
accessToken,
pkiSync.destinationConfig.vaultBaseUrl
);
logger.info(
{
syncId: pkiSync.id,
existingCertCount: Object.keys(vaultCertificates).length,
disabledCertCount: disabledAzureKeyVaultCertificateKeys.length
},
"Retrieved existing certificates from Azure Key Vault"
destinationConfig.vaultBaseUrl,
pkiSync.id
);
const setCertificates: {
@@ -150,10 +352,6 @@ export const azureKeyVaultPkiSyncFactory = ({ kmsService, appConnectionDAL }: TA
// Iterate through certificates to sync to Azure Key Vault
Object.entries(certificateMap).forEach(([certName, { cert, privateKey }]) => {
if (disabledAzureKeyVaultCertificateKeys.includes(certName)) {
logger.debug(
{ syncId: pkiSync.id, certificateName: certName },
"Skipping disabled certificate in Azure Key Vault"
);
return;
}
@@ -166,124 +364,107 @@ export const azureKeyVaultPkiSyncFactory = ({ kmsService, appConnectionDAL }: TA
cert,
privateKey
});
logger.debug(
{ syncId: pkiSync.id, certificateName: certName, isUpdate: !!existingCert },
"Certificate will be uploaded to Azure Key Vault"
);
} else {
logger.debug(
{ syncId: pkiSync.id, certificateName: certName },
"Certificate already up to date in Azure Key Vault"
);
}
});
// Identify expired/removed certificates that need to be cleaned up from Azure Key Vault
// Only remove certificates that were managed by Infisical (start with 'Infisical-')
// Only remove certificates that were managed by Infisical (match naming schema)
const certificatesToRemove = Object.keys(vaultCertificates).filter(
(vaultCertName) =>
vaultCertName.startsWith("Infisical-") &&
isInfisicalManagedCertificate(vaultCertName, pkiSync) &&
!activeCertificateNames.includes(vaultCertName) &&
!disabledAzureKeyVaultCertificateKeys.includes(vaultCertName)
);
logger.info(
{
syncId: pkiSync.id,
certificatesToUpload: setCertificates.length,
certificatesToRemove: certificatesToRemove.length,
totalCertificates: Object.keys(certificateMap).length
},
"Determined certificates to upload and remove from Azure Key Vault"
);
// Upload certificates to Azure Key Vault with rate limiting
const uploadResults = await executeWithConcurrencyLimit(
setCertificates,
async ({ key, cert, privateKey }) => {
try {
// Combine certificate and private key in PEM format for Azure Key Vault
// Azure Key Vault accepts PEM format with both cert and private key
let combinedPem = cert;
if (privateKey) {
combinedPem = `${privateKey}\n${cert}`;
}
// Upload certificates to Azure Key Vault
const uploadPromises = setCertificates.map(async ({ key, cert, privateKey }) => {
try {
// Combine certificate and private key in PEM format for Azure Key Vault
// Azure Key Vault accepts PEM format with both cert and private key
let combinedPem = cert;
if (privateKey) {
combinedPem = `${privateKey}\n${cert}`;
}
// Convert to base64 for Azure Key Vault import
const base64Cert = Buffer.from(combinedPem).toString("base64");
// Convert to base64 for Azure Key Vault import
const base64Cert = Buffer.from(combinedPem).toString("base64");
// Parse certificate to extract X.509 properties and key properties
const x509Props = parseCertificateX509Props(cert);
const keyProps = parseCertificateKeyProps(cert);
const importData = {
value: base64Cert,
policy: {
key_props: {
exportable: true,
key_size: 2048,
kty: "RSA",
reuse_key: false
// Build key_props based on key type
const keyPropsConfig = {
exportable: true,
reuse_key: false,
...keyProps
};
const importData = {
value: base64Cert,
policy: {
key_props: keyPropsConfig,
secret_props: {
contentType: "application/x-pem-file"
},
x509_props: x509Props
},
secret_props: {
contentType: "application/x-pem-file"
},
x509_props: {
subject: "",
sans: {
dns_names: [],
emails: [],
upns: []
attributes: {
enabled: true,
exportable: true
}
};
const response = await request.post(
`${destinationConfig.vaultBaseUrl}/certificates/${encodeURIComponent(key)}/import?api-version=7.4`,
importData,
{
headers: {
Authorization: `Bearer ${accessToken}`,
"Content-Type": "application/json"
}
}
}
};
);
const response = await request.post(
`${pkiSync.destinationConfig.vaultBaseUrl}/certificates/${encodeURIComponent(key)}/import?api-version=7.4`,
importData,
{
headers: {
Authorization: `Bearer ${accessToken}`,
"Content-Type": "application/json"
return { key, success: true, response: response.data as unknown };
} catch (error) {
if (error instanceof AxiosError) {
const errorMessage =
error.response?.data && typeof error.response.data === "object" && "error" in error.response.data
? (error.response.data as { error?: { message?: string } }).error?.message || error.message
: error.message;
// Check if the error is due to certificate in deleted but recoverable state
const isDeletedButRecoverable =
errorMessage.includes("deleted but recoverable state") || errorMessage.includes("name cannot be reused");
if (isDeletedButRecoverable) {
logger.warn(
{ certificateKey: key, syncId: pkiSync.id },
"Certificate exists in deleted but recoverable state in Azure Key Vault - skipping upload"
);
return { key, success: false, skipped: true, reason: "Certificate in deleted but recoverable state" };
}
throw new PkiSyncError({
message: `Failed to upload certificate ${key} to Azure Key Vault: ${errorMessage}`,
cause: error,
context: {
certificateKey: key,
statusCode: error.response?.status,
responseData: error.response?.data
}
});
}
);
logger.info(
{ syncId: pkiSync.id, certificateName: key },
"Successfully uploaded certificate to Azure Key Vault"
);
return { key, success: true, response: response.data as unknown };
} catch (error) {
if (error instanceof AxiosError) {
const errorMessage =
error.response?.data && typeof error.response.data === "object" && "error" in error.response.data
? (error.response.data as { error?: { message?: string } }).error?.message || error.message
: error.message;
// Check if the error is due to certificate in deleted but recoverable state
const isDeletedButRecoverable =
errorMessage.includes("deleted but recoverable state") || errorMessage.includes("name cannot be reused");
if (isDeletedButRecoverable) {
logger.warn(
{ certificateKey: key, syncId: pkiSync.id },
"Certificate exists in deleted but recoverable state in Azure Key Vault - skipping upload"
);
return { key, success: false, skipped: true, reason: "Certificate in deleted but recoverable state" };
}
throw new PkiSyncError({
message: `Failed to upload certificate ${key} to Azure Key Vault: ${errorMessage}`,
cause: error,
context: {
certificateKey: key,
statusCode: error.response?.status,
responseData: error.response?.data
}
});
throw error;
}
throw error;
}
});
},
{ operation: "upload-certificates", syncId: pkiSync.id }
);
const results = await Promise.allSettled(uploadPromises);
const results = uploadResults;
const failedUploads = results.filter((result) => result.status === "rejected");
const fulfilledResults = results.filter((result) => result.status === "fulfilled");
@@ -298,52 +479,37 @@ export const azureKeyVaultPkiSyncFactory = ({ kmsService, appConnectionDAL }: TA
let failedRemovals = 0;
if (certificatesToRemove.length > 0) {
logger.info(
{
syncId: pkiSync.id,
certificatesToRemove: certificatesToRemove.length
},
"Removing expired/removed certificates from Azure Key Vault"
);
const removePromises = certificatesToRemove.map(async (certName) => {
try {
await request.delete(
`${pkiSync.destinationConfig.vaultBaseUrl}/certificates/${encodeURIComponent(certName)}?api-version=7.4`,
{
headers: {
Authorization: `Bearer ${accessToken}`
const removeResults = await executeWithConcurrencyLimit(
certificatesToRemove,
async (certName) => {
try {
await request.delete(
`${destinationConfig.vaultBaseUrl}/certificates/${encodeURIComponent(certName)}?api-version=7.4`,
{
headers: {
Authorization: `Bearer ${accessToken}`
}
}
}
);
logger.info(
{ syncId: pkiSync.id, certificateName: certName },
"Successfully removed expired/removed certificate from Azure Key Vault"
);
return { key: certName, success: true };
} catch (error) {
// If certificate doesn't exist (404), consider it as successfully removed
if (error instanceof AxiosError && error.response?.status === 404) {
logger.info(
{ syncId: pkiSync.id, certificateName: certName },
"Certificate not found in Azure Key Vault during sync cleanup - considering removal successful"
);
return { key: certName, success: true, alreadyRemoved: true };
return { key: certName, success: true };
} catch (error) {
// If certificate doesn't exist (404), consider it as successfully removed
if (error instanceof AxiosError && error.response?.status === 404) {
return { key: certName, success: true, alreadyRemoved: true };
}
logger.error(
{ error, syncId: pkiSync.id, certificateName: certName },
"Failed to remove expired/removed certificate from Azure Key Vault"
);
// Don't throw here - we want to continue with other operations
return { key: certName, success: false, error: error as Error };
}
logger.error(
{ error, syncId: pkiSync.id, certificateName: certName },
"Failed to remove expired/removed certificate from Azure Key Vault"
);
// Don't throw here - we want to continue with other operations
return { key: certName, success: false, error: error as Error };
}
});
const removeResults = await Promise.allSettled(removePromises);
},
{ operation: "remove-certificates", syncId: pkiSync.id }
);
const successfulRemovals = removeResults.filter(
(result) => result.status === "fulfilled" && result.value.success
);
@@ -362,136 +528,124 @@ export const azureKeyVaultPkiSyncFactory = ({ kmsService, appConnectionDAL }: TA
}
}
// Log skipped certificates for transparency
// Collect detailed information for UI feedback
const details: {
failedUploads?: Array<{ name: string; error: string }>;
failedRemovals?: Array<{ name: string; error: string }>;
skippedCertificates?: Array<{ name: string; reason: string }>;
} = {};
// Collect skipped certificate details
if (skippedUploads.length > 0) {
const skippedNames = skippedUploads.map((result) =>
result.status === "fulfilled" ? result.value.key : "unknown"
);
logger.info(
{
syncId: pkiSync.id,
skippedCertificates: skippedNames,
skippedCount: skippedUploads.length
},
"Some certificates were skipped due to Azure Key Vault constraints"
);
details.skippedCertificates = skippedUploads.map((result) => {
const certificateName = result.status === "fulfilled" ? result.value.key : "unknown";
return {
name: certificateName,
reason: "Azure Key Vault constraints or certificate already up to date"
};
});
}
logger.info(
{
syncId: pkiSync.id,
successfulUploads: successfulUploads.length,
failedUploads: failedUploads.length,
skippedUploads: skippedUploads.length,
removedCertificates,
failedRemovals,
skippedCertificates: Object.keys(certificateMap).length - setCertificates.length
},
"Azure Key Vault certificate sync completed"
);
// Collect failed upload details
if (failedUploads.length > 0) {
const failedReasons = failedUploads.map((failure) => {
details.failedUploads = failedUploads.map((failure, index) => {
const certificateName = setCertificates[index]?.key || "unknown";
let errorMessage = "Unknown error";
if (failure.status === "rejected") {
return (failure.reason as Error)?.message || "Unknown error";
errorMessage = (failure.reason as Error)?.message || "Unknown error";
}
return "Unknown error";
return {
name: certificateName,
error: errorMessage
};
});
logger.error(
{
syncId: pkiSync.id,
failedReasons,
failedUploads: details.failedUploads,
failedCount: failedUploads.length
},
"Some certificates failed to upload to Azure Key Vault"
);
throw new PkiSyncError({
message: `Failed to upload ${failedUploads.length} certificate(s) to Azure Key Vault`,
context: {
failedReasons,
totalCertificates: setCertificates.length,
failedCount: failedUploads.length
}
});
}
return {
uploaded: setCertificates.length,
removed: removedCertificates,
failedRemovals,
skipped: Object.keys(certificateMap).length - setCertificates.length
};
};
// Collect failed removal details
if (failedRemovals > 0) {
const failedRemovalNames = certificatesToRemove.slice(-failedRemovals);
details.failedRemovals = failedRemovalNames.map((certName) => ({
name: certName,
error: "Failed to remove from Azure Key Vault"
}));
const importCertificates = async (pkiSync: TAzureKeyVaultPkiSyncWithCredentials): Promise<TCertificateMap> => {
const { accessToken } = await getAzureConnectionAccessToken(pkiSync.connection.id, appConnectionDAL, kmsService);
const { vaultCertificates } = await $getAzureKeyVaultCertificates(
accessToken,
pkiSync.destinationConfig.vaultBaseUrl
);
return vaultCertificates;
};
const removeCertificates = async (pkiSync: TAzureKeyVaultPkiSyncWithCredentials, certificateNames: string[]) => {
const { accessToken } = await getAzureConnectionAccessToken(pkiSync.connection.id, appConnectionDAL, kmsService);
// Only remove certificates that are managed by Infisical (start with 'Infisical-' prefix)
const infisicalManagedCertNames = certificateNames.filter((certName) => certName.startsWith("Infisical-"));
if (infisicalManagedCertNames.length < certificateNames.length) {
logger.debug(
logger.warn(
{
syncId: pkiSync.id,
totalRequested: certificateNames.length,
infisicalManaged: infisicalManagedCertNames.length,
skipped: certificateNames.length - infisicalManagedCertNames.length
failedRemovals: details.failedRemovals,
successfulRemovals: removedCertificates
},
"Filtered out non-Infisical certificates from removal request"
"Some expired/removed certificates could not be removed from Azure Key Vault"
);
}
const removePromises = infisicalManagedCertNames.map(async (certName) => {
try {
const response = await request.delete(
`${pkiSync.destinationConfig.vaultBaseUrl}/certificates/${encodeURIComponent(certName)}?api-version=7.4`,
{
headers: {
Authorization: `Bearer ${accessToken}`
}
}
);
return {
uploaded: successfulUploads.length,
removed: removedCertificates,
failedRemovals,
skipped: Object.keys(certificateMap).length - setCertificates.length,
details: Object.keys(details).length > 0 ? details : undefined
};
};
return { key: certName, success: true, response: response.data as unknown };
} catch (error) {
if (error instanceof AxiosError) {
// If certificate doesn't exist (404), consider it as successfully removed
if (error.response?.status === 404) {
logger.info(
{ syncId: pkiSync.id, certificateName: certName },
"Certificate not found in Azure Key Vault - considering removal successful"
);
return { key: certName, success: true, alreadyRemoved: true };
}
const removeCertificates = async (pkiSync: TPkiSyncWithCredentials, certificateNames: string[]) => {
const { accessToken } = await getAzureConnectionAccessToken(pkiSync.connection.id, appConnectionDAL, kmsService);
throw new PkiSyncError({
message: `Failed to remove certificate ${certName} from Azure Key Vault`,
cause: error,
context: {
certificateKey: certName,
statusCode: error.response?.status,
responseData: error.response?.data
// Cast destination config to Azure Key Vault config
const destinationConfig = pkiSync.destinationConfig as TAzureKeyVaultPkiSyncConfig;
// Only remove certificates that are managed by Infisical (match naming schema)
const infisicalManagedCertNames = certificateNames.filter((certName) =>
isInfisicalManagedCertificate(certName, pkiSync)
);
const results = await executeWithConcurrencyLimit(
infisicalManagedCertNames,
async (certName) => {
try {
const response = await request.delete(
`${destinationConfig.vaultBaseUrl}/certificates/${encodeURIComponent(certName)}?api-version=7.4`,
{
headers: {
Authorization: `Bearer ${accessToken}`
}
}
});
);
return { key: certName, success: true, response: response.data as unknown };
} catch (error) {
if (error instanceof AxiosError) {
// If certificate doesn't exist (404), consider it as successfully removed
if (error.response?.status === 404) {
return { key: certName, success: true, alreadyRemoved: true };
}
throw new PkiSyncError({
message: `Failed to remove certificate ${certName} from Azure Key Vault`,
cause: error,
context: {
certificateKey: certName,
statusCode: error.response?.status,
responseData: error.response?.data
}
});
}
throw error;
}
throw error;
}
});
const results = await Promise.allSettled(removePromises);
},
{ operation: "remove-specific-certificates", syncId: pkiSync.id }
);
const failedRemovals = results.filter((result) => result.status === "rejected");
if (failedRemovals.length > 0) {
@@ -521,7 +675,6 @@ export const azureKeyVaultPkiSyncFactory = ({ kmsService, appConnectionDAL }: TA
return {
syncCertificates,
importCertificates,
removeCertificates
};
};

View File

@@ -1,14 +1,45 @@
import RE2 from "re2";
import { z } from "zod";
import { AppConnection } from "@app/services/app-connection/app-connection-enums";
import { PkiSync } from "@app/services/pki-sync/pki-sync-enums";
import { PkiSyncSchema } from "@app/services/pki-sync/pki-sync-schemas";
import { AzureKeyVaultPkiSyncConfigSchema } from "./azure-key-vault-pki-sync-types";
export const AzureKeyVaultPkiSyncConfigSchema = z.object({
vaultBaseUrl: z.string().url()
});
const AzureKeyVaultPkiSyncOptionsSchema = z.object({
canImportCertificates: z.boolean().default(false),
canRemoveCertificates: z.boolean().default(true),
certificateNameSchema: z
.string()
.optional()
.refine(
(schema) => {
if (!schema) return true;
const testName = schema
.replace(new RE2("\\{\\{certificateId\\}\\}", "g"), "")
.replace(new RE2("\\{\\{environment\\}\\}", "g"), "");
const azureNamePattern = new RE2("^[a-zA-Z0-9-]{1,127}$");
const forbiddenChars = "!@#$%^&*()+=[]{}|\\:;\"'<>,.?/~` _";
const hasForbiddenChars = forbiddenChars.split("").some((char) => testName.includes(char));
return azureNamePattern.test(testName) && !hasForbiddenChars;
},
{
message:
"Certificate name schema must result in names that contain only alphanumeric characters and hyphens (a-z, A-Z, 0-9, -) and be 1-127 characters long when compiled for Azure Key Vault"
}
)
});
export const AzureKeyVaultPkiSyncSchema = PkiSyncSchema.extend({
destination: z.literal(PkiSync.AzureKeyVault),
destinationConfig: AzureKeyVaultPkiSyncConfigSchema
destinationConfig: AzureKeyVaultPkiSyncConfigSchema,
syncOptions: AzureKeyVaultPkiSyncOptionsSchema
});
export const CreateAzureKeyVaultPkiSyncSchema = z.object({
@@ -16,7 +47,7 @@ export const CreateAzureKeyVaultPkiSyncSchema = z.object({
description: z.string().optional(),
isAutoSyncEnabled: z.boolean().default(true),
destinationConfig: AzureKeyVaultPkiSyncConfigSchema,
syncOptions: z.record(z.unknown()).optional().default({}),
syncOptions: AzureKeyVaultPkiSyncOptionsSchema.optional().default({}),
subscriberId: z.string().optional(),
connectionId: z.string(),
projectId: z.string().trim().min(1)
@@ -27,7 +58,7 @@ export const UpdateAzureKeyVaultPkiSyncSchema = z.object({
description: z.string().optional(),
isAutoSyncEnabled: z.boolean().optional(),
destinationConfig: AzureKeyVaultPkiSyncConfigSchema.optional(),
syncOptions: z.record(z.unknown()).optional(),
syncOptions: AzureKeyVaultPkiSyncOptionsSchema.optional(),
subscriberId: z.string().optional(),
connectionId: z.string().optional()
});

View File

@@ -1,6 +1,13 @@
import { z } from "zod";
import { TPkiSyncWithCredentials } from "../pki-sync-types";
import { TAzureKeyVaultConnection } from "@app/services/app-connection/azure-key-vault";
import {
AzureKeyVaultPkiSyncConfigSchema,
AzureKeyVaultPkiSyncSchema,
CreateAzureKeyVaultPkiSyncSchema,
UpdateAzureKeyVaultPkiSyncSchema
} from "./azure-key-vault-pki-sync-schemas";
export type GetAzureKeyVaultCertificate = {
id: string;
@@ -18,12 +25,14 @@ export type GetAzureKeyVaultCertificate = {
cer?: string;
};
export const AzureKeyVaultPkiSyncConfigSchema = z.object({
vaultBaseUrl: z.string().url()
});
export type TAzureKeyVaultPkiSyncConfig = z.infer<typeof AzureKeyVaultPkiSyncConfigSchema>;
export type TAzureKeyVaultPkiSyncWithCredentials = TPkiSyncWithCredentials & {
destinationConfig: TAzureKeyVaultPkiSyncConfig;
export type TAzureKeyVaultPkiSync = z.infer<typeof AzureKeyVaultPkiSyncSchema>;
export type TAzureKeyVaultPkiSyncInput = z.infer<typeof CreateAzureKeyVaultPkiSyncSchema>;
export type TAzureKeyVaultPkiSyncUpdate = z.infer<typeof UpdateAzureKeyVaultPkiSyncSchema>;
export type TAzureKeyVaultPkiSyncWithCredentials = TAzureKeyVaultPkiSync & {
connection: TAzureKeyVaultConnection;
};

View File

@@ -3,16 +3,10 @@ export enum PkiSync {
}
export enum PkiSyncStatus {
Pending = "PENDING",
Running = "RUNNING",
Success = "SUCCESS",
Failed = "FAILED"
}
export enum PkiSyncImportBehavior {
ImportAllSecrets = "IMPORT_ALL_SECRETS",
PreferInfisicalSecrets = "PREFER_INFISICAL_SECRETS",
PreferExternalSecrets = "PREFER_EXTERNAL_SECRETS"
Pending = "pending",
Running = "running",
Succeeded = "succeeded",
Failed = "failed"
}
export enum PkiSyncAction {

View File

@@ -1,11 +1,16 @@
import * as handlebars from "handlebars";
import { z, ZodSchema } from "zod";
import { TLicenseServiceFactory } from "@app/ee/services/license/license-service";
import { BadRequestError } from "@app/lib/errors";
import { logger } from "@app/lib/logger";
import { TAppConnectionDALFactory } from "@app/services/app-connection/app-connection-dal";
import { TKmsServiceFactory } from "@app/services/kms/kms-service";
import { AZURE_KEY_VAULT_PKI_SYNC_LIST_OPTION } from "./azure-key-vault/azure-key-vault-pki-sync-fns";
import {
AZURE_KEY_VAULT_PKI_SYNC_LIST_OPTION,
azureKeyVaultPkiSyncFactory
} from "./azure-key-vault/azure-key-vault-pki-sync-fns";
import { PkiSync } from "./pki-sync-enums";
import { TCertificateMap, TPkiSyncWithCredentials } from "./pki-sync-types";
@@ -34,6 +39,18 @@ export const listPkiSyncOptions = () => {
return Object.values(PKI_SYNC_LIST_OPTIONS).sort((a, b) => a.name.localeCompare(b.name));
};
export const getPkiSyncProviderCapabilities = (destination: PkiSync) => {
const providerOption = PKI_SYNC_LIST_OPTIONS[destination];
if (!providerOption) {
throw new BadRequestError({ message: `Unsupported PKI sync destination: ${destination}` });
}
return {
canImportCertificates: providerOption.canImportCertificates,
canRemoveCertificates: providerOption.canRemoveCertificates
};
};
export const matchesSchema = <T extends ZodSchema>(schema: T, data: unknown): data is z.infer<T> => {
return schema.safeParse(data).success;
};
@@ -50,9 +67,124 @@ export const parsePkiSyncErrorMessage = (error: unknown): string => {
return "An unknown error occurred during PKI sync operation";
};
export const applyCertificateNameSchema = (
certificateMap: TCertificateMap,
environment: string,
schema?: string
): TCertificateMap => {
if (!schema) return certificateMap;
const processedCertificateMap: TCertificateMap = {};
for (const [certificateId, value] of Object.entries(certificateMap)) {
const newName = handlebars.compile(schema)({
certificateId,
environment
});
processedCertificateMap[newName] = value;
}
return processedCertificateMap;
};
export const stripCertificateNameSchema = (
certificateMap: TCertificateMap,
environment: string,
schema?: string
): TCertificateMap => {
if (!schema) return certificateMap;
const compiledSchemaPattern = handlebars.compile(schema)({
certificateId: "{{certificateId}}",
environment
});
const parts = compiledSchemaPattern.split("{{certificateId}}");
const prefix = parts[0];
const suffix = parts[parts.length - 1];
const strippedMap: TCertificateMap = {};
for (const [name, value] of Object.entries(certificateMap)) {
if (!name.startsWith(prefix) || !name.endsWith(suffix)) {
// eslint-disable-next-line no-continue
continue;
}
const strippedName = name.slice(prefix.length, name.length - suffix.length);
strippedMap[strippedName] = value;
}
return strippedMap;
};
export const matchesCertificateNameSchema = (name: string, environment: string, schema?: string): boolean => {
if (!schema) return true;
const compiledSchemaPattern = handlebars.compile(schema)({
certificateId: "{{certificateId}}",
environment
});
if (!compiledSchemaPattern.includes("{{certificateId}}")) {
return name === compiledSchemaPattern;
}
const parts = compiledSchemaPattern.split("{{certificateId}}");
const prefix = parts[0];
const suffix = parts[parts.length - 1];
if (prefix === "" && suffix === "") return true;
// If prefix is empty, name must end with suffix
if (prefix === "") return name.endsWith(suffix);
// If suffix is empty, name must start with prefix
if (suffix === "") return name.startsWith(prefix);
// Name must start with prefix and end with suffix
return name.startsWith(prefix) && name.endsWith(suffix);
};
const isAzureKeyVaultPkiSync = (pkiSync: TPkiSyncWithCredentials): boolean => {
return pkiSync.destination === PkiSync.AzureKeyVault;
};
/**
* Trigger auto sync for PKI syncs connected to a PKI subscriber when certificates are issued/revoked/deleted
*/
export const triggerAutoSyncForSubscriber = async (
subscriberId: string,
dependencies: {
pkiSyncDAL: {
find: (filter: { subscriberId: string; isAutoSyncEnabled: boolean }) => Promise<Array<{ id: string }>>;
};
pkiSyncQueue: {
queuePkiSyncSyncCertificatesById: (params: { syncId: string }) => Promise<void>;
};
}
) => {
try {
const pkiSyncs = await dependencies.pkiSyncDAL.find({
subscriberId,
isAutoSyncEnabled: true
});
// Queue sync jobs for each auto sync enabled PKI sync
const syncPromises = pkiSyncs.map((pkiSync) =>
dependencies.pkiSyncQueue.queuePkiSyncSyncCertificatesById({ syncId: pkiSync.id })
);
await Promise.all(syncPromises);
} catch (error) {
logger.error(error, `Failed to trigger auto sync for subscriber ${subscriberId}:`);
}
};
export const PkiSyncFns = {
getCertificates: async (
pkiSync: TPkiSyncWithCredentials,
// eslint-disable-next-line @typescript-eslint/no-unused-vars
dependencies: {
appConnectionDAL: Pick<TAppConnectionDALFactory, "findById" | "updateById">;
kmsService: Pick<TKmsServiceFactory, "createCipherPairWithDataKey">;
@@ -60,11 +192,8 @@ export const PkiSyncFns = {
): Promise<TCertificateMap> => {
switch (pkiSync.destination) {
case PkiSync.AzureKeyVault: {
const { azureKeyVaultPkiSyncFactory } = await import("./azure-key-vault/azure-key-vault-pki-sync-fns");
const azureKeyVaultPkiSync = azureKeyVaultPkiSyncFactory(dependencies);
// Type assertion needed due to destinationConfig type differences
return azureKeyVaultPkiSync.importCertificates(
pkiSync as unknown as import("./azure-key-vault/azure-key-vault-pki-sync-types").TAzureKeyVaultPkiSyncWithCredentials
throw new Error(
"Azure Key Vault does not support importing certificates into Infisical (private keys cannot be extracted)"
);
}
default:
@@ -84,16 +213,19 @@ export const PkiSyncFns = {
removed?: number;
failedRemovals?: number;
skipped: number;
details?: {
failedUploads?: Array<{ name: string; error: string }>;
failedRemovals?: Array<{ name: string; error: string }>;
skippedCertificates?: Array<{ name: string; reason: string }>;
};
}> => {
switch (pkiSync.destination) {
case PkiSync.AzureKeyVault: {
const { azureKeyVaultPkiSyncFactory } = await import("./azure-key-vault/azure-key-vault-pki-sync-fns");
if (!isAzureKeyVaultPkiSync(pkiSync)) {
throw new Error("Invalid Azure Key Vault PKI sync configuration");
}
const azureKeyVaultPkiSync = azureKeyVaultPkiSyncFactory(dependencies);
// Type assertion needed due to destinationConfig type differences
return azureKeyVaultPkiSync.syncCertificates(
pkiSync as unknown as import("./azure-key-vault/azure-key-vault-pki-sync-types").TAzureKeyVaultPkiSyncWithCredentials,
certificateMap
);
return azureKeyVaultPkiSync.syncCertificates(pkiSync, certificateMap);
}
default:
throw new Error(`Unsupported PKI sync destination: ${String(pkiSync.destination)}`);
@@ -110,13 +242,11 @@ export const PkiSyncFns = {
): Promise<void> => {
switch (pkiSync.destination) {
case PkiSync.AzureKeyVault: {
const { azureKeyVaultPkiSyncFactory } = await import("./azure-key-vault/azure-key-vault-pki-sync-fns");
if (!isAzureKeyVaultPkiSync(pkiSync)) {
throw new Error("Invalid Azure Key Vault PKI sync configuration");
}
const azureKeyVaultPkiSync = azureKeyVaultPkiSyncFactory(dependencies);
// Type assertion needed due to destinationConfig type differences
await azureKeyVaultPkiSync.removeCertificates(
pkiSync as unknown as import("./azure-key-vault/azure-key-vault-pki-sync-types").TAzureKeyVaultPkiSyncWithCredentials,
certificateNames
);
await azureKeyVaultPkiSync.removeCertificates(pkiSync, certificateNames);
break;
}
default:

View File

@@ -3,8 +3,8 @@ import opentelemetry from "@opentelemetry/api";
import * as x509 from "@peculiar/x509";
import { AxiosError } from "axios";
import { Job } from "bullmq";
import handlebars from "handlebars";
import { ProjectMembershipRole } from "@app/db/schemas";
import { EventType, TAuditLogServiceFactory } from "@app/ee/services/audit-log/audit-log-types";
import { TLicenseServiceFactory } from "@app/ee/services/license/license-service";
import { KeyStorePrefixes, TKeyStoreFactory } from "@app/keystore/keystore";
@@ -16,20 +16,17 @@ import { ActorType } from "@app/services/auth/auth-type";
import { TKmsServiceFactory } from "@app/services/kms/kms-service";
import { TProjectDALFactory } from "@app/services/project/project-dal";
import { getProjectKmsCertificateKeyId } from "@app/services/project/project-fns";
import { TProjectMembershipDALFactory } from "@app/services/project-membership/project-membership-dal";
import { TAppConnectionDALFactory } from "../app-connection/app-connection-dal";
import { TCertificateBodyDALFactory } from "../certificate/certificate-body-dal";
import { TCertificateDALFactory } from "../certificate/certificate-dal";
import { getCertificateCredentials } from "../certificate/certificate-fns";
import { TCertificateSecretDALFactory } from "../certificate/certificate-secret-dal";
import { CertStatus } from "../certificate/certificate-types";
import { TPkiSyncDALFactory } from "./pki-sync-dal";
import { PkiSyncAction } from "./pki-sync-enums";
import { PkiSyncStatus } from "./pki-sync-enums";
import { PkiSyncError } from "./pki-sync-errors";
import { enterprisePkiSyncCheck, parsePkiSyncErrorMessage, PkiSyncFns } from "./pki-sync-fns";
import {
PkiSyncStatus,
TCertificateMap,
TPkiSyncImportCertificatesDTO,
TPkiSyncRaw,
@@ -38,9 +35,7 @@ import {
TPkiSyncWithCredentials,
TQueuePkiSyncImportCertificatesByIdDTO,
TQueuePkiSyncRemoveCertificatesByIdDTO,
TQueuePkiSyncSyncCertificatesByIdDTO,
TQueueSendPkiSyncActionFailedNotificationsDTO,
TSendPkiSyncFailedNotificationsJobDTO
TQueuePkiSyncSyncCertificatesByIdDTO
} from "./pki-sync-types";
export type TPkiSyncQueueFactory = ReturnType<typeof pkiSyncQueueFactory>;
@@ -55,7 +50,6 @@ type TPkiSyncQueueFactoryDep = {
keyStore: Pick<TKeyStoreFactory, "acquireLock" | "setItemWithExpiry" | "getItem">;
pkiSyncDAL: Pick<TPkiSyncDALFactory, "findById" | "find" | "updateById" | "deleteById" | "update">;
auditLogService: Pick<TAuditLogServiceFactory, "createAuditLog">;
projectMembershipDAL: Pick<TProjectMembershipDALFactory, "findAllProjectMembers">;
projectDAL: TProjectDALFactory;
licenseService: Pick<TLicenseServiceFactory, "getPlan">;
certificateDAL: Pick<
@@ -88,7 +82,6 @@ export const pkiSyncQueueFactory = ({
keyStore,
pkiSyncDAL,
auditLogService,
projectMembershipDAL,
projectDAL,
licenseService,
certificateDAL,
@@ -151,89 +144,6 @@ export const pkiSyncQueueFactory = ({
);
};
const $createCertificatesInSubscriber = async (
pkiSync: TPkiSyncWithCredentials,
certificatesToCreate: Array<{
name: string;
certificate: string;
privateKey?: string;
}>
) => {
const { projectId, subscriberId } = pkiSync;
if (!subscriberId) {
throw new Error("PKI Sync subscriber ID is required for certificate creation");
}
logger.info(`Creating ${certificatesToCreate.length} certificates in PKI subscriber ${subscriberId}`);
for (const certData of certificatesToCreate) {
try {
// Validate certificate data
if (!certData.certificate || certData.certificate.trim() === "") {
logger.error(`Skipping certificate ${certData.name}: empty certificate data`);
// eslint-disable-next-line no-continue
continue;
}
// Parse certificate to extract metadata
const cert = new x509.X509Certificate(certData.certificate);
const { serialNumber } = cert;
const { notBefore } = cert;
const { notAfter } = cert;
const commonName =
cert.subject
.split(",")
.find((part) => part.trim().startsWith("CN="))
?.split("=")[1]
?.trim() || certData.name;
// Get KMS key for encryption
const kmsKeyId = await getProjectKmsCertificateKeyId({ projectId, projectDAL, kmsService });
const kmsEncryptor = await kmsService.encryptWithKmsKey({
kmsId: kmsKeyId
});
// Create certificate record
const createdCert = await certificateDAL.create({
pkiSubscriberId: subscriberId,
status: CertStatus.ACTIVE,
serialNumber,
notBefore,
notAfter,
commonName,
friendlyName: certData.name,
projectId
});
// Create certificate body record with encrypted certificate
const { cipherTextBlob: encryptedCertificate } = await kmsEncryptor({
plainText: Buffer.from(certData.certificate)
});
await certificateBodyDAL.create({
certId: createdCert.id,
encryptedCertificate
});
if (certData.privateKey) {
const { cipherTextBlob: encryptedPrivateKey } = await kmsEncryptor({
plainText: Buffer.from(certData.privateKey)
});
await certificateSecretDAL.create({
certId: createdCert.id,
encryptedPrivateKey
});
}
logger.info(`Successfully created certificate ${certData.name} with ID ${createdCert.id}`);
} catch (error) {
logger.error(`Failed to create certificate ${certData.name}: ${String(error)}`);
}
}
};
const $getInfisicalCertificates = async (
pkiSync: TPkiSyncRaw | TPkiSyncWithCredentials
): Promise<TCertificateMap> => {
@@ -254,34 +164,8 @@ export const pkiSyncQueueFactory = ({
subscriberId
});
logger.info(
{ subscriberId, certificateCount: certificates.length },
"Found active certificates for PKI sync subscriber"
);
for (const certificate of certificates) {
try {
// Only sync certificates issued by Infisical (not imported ones)
if (!certificate.caId) {
logger.debug(
{ certificateId: certificate.id, subscriberId },
"Skipping imported certificate - not syncing to destination"
);
// eslint-disable-next-line no-continue
continue;
}
// Check if certificate is expired
const now = new Date();
if (certificate.notAfter < now) {
logger.debug(
{ certificateId: certificate.id, subscriberId, expiredAt: certificate.notAfter },
"Skipping expired certificate"
);
// eslint-disable-next-line no-continue
continue;
}
// Get the certificate body and decrypt the certificate data
const certBody = await certificateBodyDAL.findOne({ certId: certificate.id });
@@ -323,19 +207,24 @@ export const pkiSyncQueueFactory = ({
certPrivateKey = undefined;
}
// Use Infisical-prefixed ID for clear identification in destination
// Azure Key Vault doesn't allow underscores, so use hyphens and remove UUID hyphens
const certificateName = `Infisical-${certificate.id.replace(/-/g, "")}`;
let certificateName: string;
const syncOptions = pkiSync.syncOptions as { certificateNameSchema?: string } | undefined;
const certificateNameSchema = syncOptions?.certificateNameSchema;
if (certificateNameSchema) {
const environment = "global";
certificateName = handlebars.compile(certificateNameSchema)({
certificateId: certificate.id.replace(/-/g, ""),
environment
});
} else {
certificateName = `Infisical-${certificate.id.replace(/-/g, "")}`;
}
certificateMap[certificateName] = {
cert: certificatePem,
privateKey: certPrivateKey || ""
};
logger.info(
{ certificateId: certificate.id, certificateName, subscriberId },
"Successfully prepared certificate for PKI sync"
);
} else {
logger.warn({ certificateId: certificate.id, subscriberId }, "Certificate body not found for certificate");
}
@@ -395,83 +284,8 @@ export const pkiSyncQueueFactory = ({
removeOnFail: true
});
const $queueSendPkiSyncFailedNotifications = async (payload: TQueueSendPkiSyncActionFailedNotificationsDTO) => {
if (!appCfg.isSmtpConfigured) return;
await queueService.queue(QueueName.PkiSync, QueueJobs.PkiSyncSendActionFailedNotifications, payload, {
jobId: `pki-sync-${payload.pkiSync.id}-failed-notifications`,
attempts: 5,
delay: 1000 * 60,
backoff: {
type: "exponential",
delay: 3000
},
removeOnFail: true,
removeOnComplete: true
});
};
const $importCertificates = async (pkiSync: TPkiSyncWithCredentials): Promise<TCertificateMap> => {
const {
projectId,
destination,
connection: { orgId }
} = pkiSync;
await enterprisePkiSyncCheck(
licenseService,
orgId,
destination,
"Failed to import certificates due to plan restriction. Upgrade plan to access enterprise PKI syncs."
);
if (!projectId) {
throw new Error("Invalid PKI Sync source configuration: project no longer exists.");
}
const importedCertificates = await PkiSyncFns.getCertificates(pkiSync, {
appConnectionDAL,
kmsService
});
if (!Object.keys(importedCertificates).length) return {};
const importedCertificateMap: TCertificateMap = {};
const certificateMap = await $getInfisicalCertificates(pkiSync);
// Compare existing certificates with imported ones and determine which need to be created/updated
const certificatesToCreate: Array<{
name: string;
certificate: string;
privateKey?: string;
}> = [];
Object.entries(importedCertificates).forEach(([name, certificateData]) => {
const { cert: certificate, privateKey } = certificateData;
if (!Object.prototype.hasOwnProperty.call(certificateMap, name)) {
// Certificate doesn't exist in Infisical, create it
certificatesToCreate.push({
name,
certificate,
privateKey
});
importedCertificateMap[name] = certificateData;
} else {
// Certificate exists - could compare and update if needed
// For now, we'll skip updating existing certificates to avoid conflicts
importedCertificateMap[name] = certificateData;
}
});
// Create new certificates in Infisical
if (certificatesToCreate.length > 0) {
logger.info(`PKI Sync Import: Creating ${certificatesToCreate.length} new certificates`);
await $createCertificatesInSubscriber(pkiSync, certificatesToCreate);
}
return importedCertificateMap;
const $importCertificates = async (): Promise<TCertificateMap> => {
throw new Error("Certificate import functionality is not implemented");
};
const $handleSyncCertificatesJob = async (job: TPkiSyncSyncCertificatesDTO, pkiSync: TPkiSyncRaw) => {
@@ -586,24 +400,14 @@ export const pkiSyncQueueFactory = ({
});
if (isSynced || isFinalAttempt) {
const updatedPkiSync = await pkiSyncDAL.updateById(pkiSync.id, {
await pkiSyncDAL.updateById(pkiSync.id, {
syncStatus,
lastSyncJobId: job.id,
lastSyncMessage: syncMessage,
lastSyncedAt: isSynced ? ranAt : undefined
});
if (!isSynced) {
await $queueSendPkiSyncFailedNotifications({
pkiSync: updatedPkiSync,
action: PkiSyncAction.SyncCertificates,
auditLogInfo
});
}
}
}
logger.info("PkiSync Sync Job with ID %s Completed", job.id);
};
const $handleImportCertificatesJob = async (job: TPkiSyncImportCertificatesDTO, pkiSync: TPkiSyncRaw) => {
@@ -624,24 +428,7 @@ export const pkiSyncQueueFactory = ({
let isFinalAttempt = job.attemptsStarted === job.opts.attempts;
try {
const {
connection: { orgId, encryptedCredentials, projectId: appConnectionProjectId }
} = pkiSync;
const credentials = await decryptAppConnectionCredentials({
orgId,
encryptedCredentials,
kmsService,
projectId: appConnectionProjectId
});
await $importCertificates({
...pkiSync,
connection: {
...pkiSync.connection,
credentials
}
} as TPkiSyncWithCredentials);
await $importCertificates();
isSuccess = true;
} catch (err) {
@@ -693,24 +480,14 @@ export const pkiSyncQueueFactory = ({
});
if (isSuccess || isFinalAttempt) {
const updatedPkiSync = await pkiSyncDAL.updateById(pkiSync.id, {
await pkiSyncDAL.updateById(pkiSync.id, {
importStatus,
lastImportJobId: job.id,
lastImportMessage: importMessage,
lastImportedAt: isSuccess ? ranAt : undefined
});
if (!isSuccess) {
await $queueSendPkiSyncFailedNotifications({
pkiSync: updatedPkiSync,
action: PkiSyncAction.ImportCertificates,
auditLogInfo
});
}
}
}
logger.info("PkiSync Import Job with ID %s Completed", job.id);
};
const $handleRemoveCertificatesJob = async (job: TPkiSyncRemoveCertificatesDTO, pkiSync: TPkiSyncRaw) => {
@@ -819,72 +596,19 @@ export const pkiSyncQueueFactory = ({
if (isSuccess && deleteSyncOnComplete) {
await pkiSyncDAL.deleteById(pkiSync.id);
} else {
const updatedPkiSync = await pkiSyncDAL.updateById(pkiSync.id, {
await pkiSyncDAL.updateById(pkiSync.id, {
removeStatus,
lastRemoveJobId: job.id,
lastRemoveMessage: removeMessage,
lastRemovedAt: isSuccess ? ranAt : undefined
});
if (!isSuccess) {
await $queueSendPkiSyncFailedNotifications({
pkiSync: updatedPkiSync,
action: PkiSyncAction.RemoveCertificates,
auditLogInfo
});
}
}
}
}
logger.info("PkiSync Remove Job with ID %s Completed", job.id);
};
const $sendPkiSyncFailedNotifications = async (job: TSendPkiSyncFailedNotificationsJobDTO) => {
const {
data: { pkiSync, auditLogInfo, action }
} = job;
const { projectId, name, lastSyncMessage, lastRemoveMessage, lastImportMessage } = pkiSync;
const projectMembers = await projectMembershipDAL.findAllProjectMembers(projectId);
const project = await projectDAL.findById(projectId);
// Filter for project admins similar to secret sync
let projectAdmins = projectMembers.filter((member) =>
member.roles.some((role) => role.role === ProjectMembershipRole.Admin)
);
const triggeredByUserId = auditLogInfo?.actor?.type === ActorType.USER ? auditLogInfo.actor.metadata?.userId : null;
if (triggeredByUserId) {
// Don't send notification to the user who triggered the action
projectAdmins = projectAdmins.filter((member) => member.user.id !== triggeredByUserId);
}
// Get appropriate error message based on action type
let errorMessage: string | null = null;
if (action === PkiSyncAction.SyncCertificates) {
errorMessage = lastSyncMessage || null;
} else if (action === PkiSyncAction.ImportCertificates) {
errorMessage = lastImportMessage || null;
} else {
errorMessage = lastRemoveMessage || null;
}
if (projectAdmins.length > 0) {
logger.info(
`PKI Sync ${action} failure notification would be sent to ${projectAdmins.length} admin(s) for sync "${name}" in project "${project.name}". Error: ${errorMessage}`
);
} else {
logger.info(
`PKI Sync ${action} failure occurred for sync "${name}" in project "${project.name}" but no admins to notify. Error: ${errorMessage}`
);
}
};
const $handleAcquireLockFailure = async (job: PkiSyncActionJob) => {
const { syncId, auditLogInfo } = job.data;
const { syncId } = job.data;
switch (job.name) {
case QueueJobs.PkiSyncSyncCertificates: {
@@ -895,51 +619,33 @@ export const pkiSyncQueueFactory = ({
return;
}
const pkiSync = await pkiSyncDAL.updateById(syncId, {
await pkiSyncDAL.updateById(syncId, {
syncStatus: PkiSyncStatus.Failed,
lastSyncMessage:
"Failed to run job. This typically happens when a sync is already in progress. Please try again.",
lastSyncJobId: job.id
});
await $queueSendPkiSyncFailedNotifications({
pkiSync,
action: PkiSyncAction.SyncCertificates,
auditLogInfo
});
break;
}
case QueueJobs.PkiSyncImportCertificates: {
const pkiSync = await pkiSyncDAL.updateById(syncId, {
await pkiSyncDAL.updateById(syncId, {
importStatus: PkiSyncStatus.Failed,
lastImportMessage:
"Failed to run job. This typically happens when a sync is already in progress. Please try again.",
lastImportJobId: job.id
});
await $queueSendPkiSyncFailedNotifications({
pkiSync,
action: PkiSyncAction.ImportCertificates,
auditLogInfo
});
break;
}
case QueueJobs.PkiSyncRemoveCertificates: {
const pkiSync = await pkiSyncDAL.updateById(syncId, {
await pkiSyncDAL.updateById(syncId, {
removeStatus: PkiSyncStatus.Failed,
lastRemoveMessage:
"Failed to run job. This typically happens when a sync is already in progress. Please try again.",
lastRemoveJobId: job.id
});
await $queueSendPkiSyncFailedNotifications({
pkiSync,
action: PkiSyncAction.RemoveCertificates,
auditLogInfo
});
break;
}
default:
@@ -949,15 +655,7 @@ export const pkiSyncQueueFactory = ({
};
queueService.start(QueueName.PkiSync, async (job) => {
if (job.name === QueueJobs.PkiSyncSendActionFailedNotifications) {
await $sendPkiSyncFailedNotifications(job as TSendPkiSyncFailedNotificationsJobDTO);
return;
}
const { syncId } = job.data as
| TQueuePkiSyncSyncCertificatesByIdDTO
| TQueuePkiSyncImportCertificatesByIdDTO
| TQueuePkiSyncRemoveCertificatesByIdDTO;
const { syncId } = job.data;
const pkiSync = await pkiSyncDAL.findById(syncId);
@@ -969,10 +667,6 @@ export const pkiSyncQueueFactory = ({
const isConcurrentLimitReached = await $isConnectionConcurrencyLimitReached(connectionId);
if (isConcurrentLimitReached) {
logger.info(
`PkiSync Concurrency limit reached [syncId=${syncId}] [job=${job.name}] [connectionId=${connectionId}]`
);
await $handleAcquireLockFailure(job as PkiSyncActionJob);
return;
@@ -988,8 +682,6 @@ export const pkiSyncQueueFactory = ({
5 * 60 * 1000
);
} catch (e) {
logger.info(`PkiSync Failed to acquire lock [syncId=${syncId}] [job=${job.name}]`);
await $handleAcquireLockFailure(job as PkiSyncActionJob);
return;

View File

@@ -1,20 +1,48 @@
import RE2 from "re2";
import { z } from "zod";
import { AzureKeyVaultPkiSyncConfigSchema } from "./azure-key-vault/azure-key-vault-pki-sync-types";
import { PkiSync } from "./pki-sync-enums";
// Schema for PKI sync options configuration
export const PkiSyncOptionsSchema = z.object({
canImportCertificates: z.boolean()
canImportCertificates: z.boolean(),
canRemoveCertificates: z.boolean().optional(),
certificateNameSchema: z
.string()
.optional()
.refine(
(val) => {
if (!val) return true;
const allowedOptionalPlaceholders = ["{{environment}}"];
const allowedPlaceholdersRegexPart = ["{{certificateId}}", ...allowedOptionalPlaceholders]
.map((p) => p.replace(/[-/\\^$*+?.()|[\]{}]/g, "\\$&")) // Escape regex special characters
.join("|");
const allowedContentRegex = new RE2(`^([a-zA-Z0-9_\\-/]|${allowedPlaceholdersRegexPart})*$`);
const contentIsValid = allowedContentRegex.test(val);
if (val.trim()) {
const certificateIdRegex = new RE2(/\{\{certificateId\}\}/);
const certificateIdIsPresent = certificateIdRegex.test(val);
return contentIsValid && certificateIdIsPresent;
}
return contentIsValid;
},
{
message:
"Certificate name schema must include exactly one {{certificateId}} placeholder. It can also include {{environment}} placeholders. Only alphanumeric characters (a-z, A-Z, 0-9), dashes (-), underscores (_), and slashes (/) are allowed besides the placeholders."
}
)
});
// Schema for destination-specific configurations
export const PkiSyncDestinationConfigSchema = z.discriminatedUnion("destination", [
z.object({
destination: z.literal(PkiSync.AzureKeyVault),
config: AzureKeyVaultPkiSyncConfigSchema
})
]);
export const PkiSyncDestinationConfigSchema = z.object({
destination: z.nativeEnum(PkiSync),
config: z.record(z.unknown())
});
// Base PKI sync schema for API responses
export const PkiSyncSchema = z.object({
@@ -33,18 +61,3 @@ export const PkiSyncSchema = z.object({
syncStatus: z.string().nullable().optional(),
lastSyncedAt: z.date().nullable().optional()
});
// Schema for PKI sync list items (includes app connection info)
export const PkiSyncListItemSchema = PkiSyncSchema.extend({
appConnectionName: z.string().max(255),
appConnectionApp: z.string().max(255)
});
export const PkiSyncDetailsSchema = PkiSyncSchema.extend({
appConnectionName: z.string().max(255),
appConnectionApp: z.string().max(255)
});
export type TPkiSyncSchema = z.infer<typeof PkiSyncSchema>;
export type TPkiSyncListItemSchema = z.infer<typeof PkiSyncListItemSchema>;
export type TPkiSyncDetailsSchema = z.infer<typeof PkiSyncDetailsSchema>;

View File

@@ -11,17 +11,15 @@ import { TAppConnectionServiceFactory } from "@app/services/app-connection/app-c
import { TPkiSubscriberDALFactory } from "@app/services/pki-subscriber/pki-subscriber-dal";
import { TPkiSyncDALFactory } from "./pki-sync-dal";
import { PkiSync } from "./pki-sync-enums";
import { enterprisePkiSyncCheck, listPkiSyncOptions } from "./pki-sync-fns";
import { PkiSync, PkiSyncStatus } from "./pki-sync-enums";
import { enterprisePkiSyncCheck, getPkiSyncProviderCapabilities, listPkiSyncOptions } from "./pki-sync-fns";
import { PKI_SYNC_CONNECTION_MAP, PKI_SYNC_NAME_MAP } from "./pki-sync-maps";
import { TPkiSyncQueueFactory } from "./pki-sync-queue";
import {
PkiSyncStatus,
TCreatePkiSyncDTO,
TDeletePkiSyncDTO,
TFindPkiSyncByIdDTO,
TFindPkiSyncByNameDTO,
TListPkiSyncsByProjectId,
TListPkiSyncsBySubscriberId,
TPkiSync,
TTriggerPkiSyncImportCertificatesByIdDTO,
TTriggerPkiSyncRemoveCertificatesByIdDTO,
@@ -30,12 +28,13 @@ import {
} from "./pki-sync-types";
const getDestinationAppType = (destination: PkiSync): AppConnection => {
switch (destination) {
case PkiSync.AzureKeyVault:
return AppConnection.AzureKeyVault;
default:
throw new BadRequestError({ message: "Unsupported PKI sync destination" });
const appConnection = PKI_SYNC_CONNECTION_MAP[destination];
if (!appConnection) {
throw new BadRequestError({
message: `Unsupported PKI sync destination: ${destination}`
});
}
return appConnection;
};
type TPkiSyncServiceFactoryDep = {
@@ -85,27 +84,30 @@ export const pkiSyncServiceFactory = ({
projectId
});
ForbiddenError.from(permission).throwUnlessCan(
ProjectPermissionPkiSyncActions.Create,
subject(ProjectPermissionSub.PkiSyncs, { projectId })
);
let subscriber;
if (subscriberId) {
const subscriber = await pkiSubscriberDAL.findById(subscriberId);
subscriber = await pkiSubscriberDAL.findById(subscriberId);
if (!subscriber || subscriber.projectId !== projectId) {
throw new NotFoundError({ message: "PKI subscriber not found" });
}
}
ForbiddenError.from(permission).throwUnlessCan(
ProjectPermissionPkiSyncActions.Create,
subscriber
? subject(ProjectPermissionSub.PkiSyncs, { subscriberName: subscriber.name })
: ProjectPermissionSub.PkiSyncs
);
// Get the destination app type based on PKI sync destination
const destinationApp = getDestinationAppType(destination);
// Validates permission to connect and app is valid for sync destination
await appConnectionService.connectAppConnectionById(destinationApp, connectionId, actor);
const defaultSyncOptions = {
canImportCertificates: false,
canRemoveCertificates: true,
const providerCapabilities = getPkiSyncProviderCapabilities(destination);
const resolvedSyncOptions = {
...providerCapabilities,
...syncOptions
};
@@ -116,7 +118,7 @@ export const pkiSyncServiceFactory = ({
destination,
isAutoSyncEnabled,
destinationConfig,
syncOptions: defaultSyncOptions,
syncOptions: resolvedSyncOptions,
subscriberId,
connectionId,
projectId,
@@ -141,7 +143,6 @@ export const pkiSyncServiceFactory = ({
const updatePkiSync = async (
{
id,
projectId,
name,
description,
isAutoSyncEnabled,
@@ -149,31 +150,38 @@ export const pkiSyncServiceFactory = ({
syncOptions,
subscriberId,
connectionId
}: Omit<TUpdatePkiSyncDTO, "auditLogInfo">,
}: Omit<TUpdatePkiSyncDTO, "auditLogInfo" | "projectId">,
actor: OrgServiceActor
): Promise<TPkiSync> => {
const existingSync = await pkiSyncDAL.findById(id);
if (!existingSync) throw new NotFoundError({ message: "PKI sync not found" });
const { permission } = await permissionService.getProjectPermission({
actor: actor.type,
actorId: actor.id,
actorAuthMethod: actor.authMethod,
actorOrgId: actor.orgId,
actionProjectType: ActionProjectType.CertificateManager,
projectId
projectId: existingSync.projectId
});
const pkiSync = await pkiSyncDAL.findByIdAndProjectId(id, projectId);
const pkiSync = await pkiSyncDAL.findByIdAndProjectId(id, existingSync.projectId);
if (!pkiSync) throw new NotFoundError({ message: "PKI sync not found" });
let currentSubscriber;
if (pkiSync.subscriberId) {
currentSubscriber = await pkiSubscriberDAL.findById(pkiSync.subscriberId);
}
ForbiddenError.from(permission).throwUnlessCan(
ProjectPermissionPkiSyncActions.Edit,
subject(ProjectPermissionSub.PkiSyncs, {
projectId,
subscriberId: pkiSync.subscriberId
})
currentSubscriber
? subject(ProjectPermissionSub.PkiSyncs, { subscriberName: currentSubscriber.name })
: ProjectPermissionSub.PkiSyncs
);
if (name && name !== pkiSync.name) {
const existingPkiSync = await pkiSyncDAL.findByNameAndProjectId(name, projectId);
const existingPkiSync = await pkiSyncDAL.findByNameAndProjectId(name, existingSync.projectId);
if (existingPkiSync) {
throw new BadRequestError({ message: "PKI sync with this name already exists" });
}
@@ -181,25 +189,44 @@ export const pkiSyncServiceFactory = ({
if (subscriberId) {
const subscriber = await pkiSubscriberDAL.findById(subscriberId);
if (!subscriber || subscriber.projectId !== projectId) {
if (!subscriber || subscriber.projectId !== existingSync.projectId) {
throw new NotFoundError({ message: "PKI subscriber not found" });
}
}
if (connectionId && connectionId !== pkiSync.connectionId) {
const destinationApp =
pkiSync.destination === PkiSync.AzureKeyVault
? AppConnection.AzureKeyVault
: (pkiSync.destination as AppConnection);
const destinationApp = getDestinationAppType(pkiSync.destination);
await appConnectionService.connectAppConnectionById(destinationApp, connectionId, actor);
}
let resolvedSyncOptions = syncOptions;
if (syncOptions) {
const providerCapabilities = getPkiSyncProviderCapabilities(pkiSync.destination);
if (syncOptions.canImportCertificates && !providerCapabilities.canImportCertificates) {
throw new BadRequestError({
message: `Certificate import is not supported for ${PKI_SYNC_NAME_MAP[pkiSync.destination]} PKI sync destination`
});
}
if (syncOptions.canRemoveCertificates === false && providerCapabilities.canRemoveCertificates) {
throw new BadRequestError({
message: `Certificate removal cannot be disabled for ${PKI_SYNC_NAME_MAP[pkiSync.destination]} PKI sync destination`
});
}
resolvedSyncOptions = {
...providerCapabilities,
...syncOptions
};
}
const updatedPkiSync = await pkiSyncDAL.updateById(id, {
name,
description,
isAutoSyncEnabled,
destinationConfig,
syncOptions,
syncOptions: resolvedSyncOptions,
subscriberId,
connectionId
});
@@ -208,31 +235,61 @@ export const pkiSyncServiceFactory = ({
};
const deletePkiSync = async (
{ id, projectId }: Omit<TDeletePkiSyncDTO, "auditLogInfo">,
{ id }: Omit<TDeletePkiSyncDTO, "auditLogInfo" | "projectId">,
actor: OrgServiceActor
): Promise<TPkiSync> => {
const existingSync = await pkiSyncDAL.findById(id);
if (!existingSync) throw new NotFoundError({ message: "PKI sync not found" });
const { permission } = await permissionService.getProjectPermission({
actor: actor.type,
actorId: actor.id,
actorAuthMethod: actor.authMethod,
actorOrgId: actor.orgId,
actionProjectType: ActionProjectType.CertificateManager,
projectId
projectId: existingSync.projectId
});
const pkiSync = await pkiSyncDAL.findByIdAndProjectId(id, projectId);
const pkiSync = await pkiSyncDAL.findByIdAndProjectId(id, existingSync.projectId);
if (!pkiSync) throw new NotFoundError({ message: "PKI sync not found" });
let pkiSyncSubscriber;
if (pkiSync.subscriberId) {
pkiSyncSubscriber = await pkiSubscriberDAL.findById(pkiSync.subscriberId);
}
ForbiddenError.from(permission).throwUnlessCan(
ProjectPermissionPkiSyncActions.Delete,
subject(ProjectPermissionSub.PkiSyncs, {
projectId,
subscriberId: pkiSync.subscriberId
})
pkiSyncSubscriber
? subject(ProjectPermissionSub.PkiSyncs, { subscriberName: pkiSyncSubscriber.name })
: ProjectPermissionSub.PkiSyncs
);
const deletedPkiSync = await pkiSyncDAL.deleteById(id);
return deletedPkiSync as TPkiSync;
await pkiSyncDAL.deleteById(id);
return {
...pkiSync,
description: pkiSync.description || undefined,
subscriberId: pkiSync.subscriberId || undefined,
syncStatus: pkiSync.syncStatus || undefined,
lastSyncedAt: pkiSync.lastSyncedAt || undefined,
lastSyncJobId: pkiSync.lastSyncJobId || undefined,
lastSyncMessage: pkiSync.lastSyncMessage || undefined,
importStatus: pkiSync.importStatus || undefined,
lastImportJobId: pkiSync.lastImportJobId || undefined,
lastImportMessage: pkiSync.lastImportMessage || undefined,
lastImportedAt: pkiSync.lastImportedAt || undefined,
removeStatus: pkiSync.removeStatus || undefined,
lastRemoveJobId: pkiSync.lastRemoveJobId || undefined,
lastRemoveMessage: pkiSync.lastRemoveMessage || undefined,
lastRemovedAt: pkiSync.lastRemovedAt || undefined,
connection: {
...pkiSync.connection,
description: pkiSync.connection.description || undefined,
gatewayId: pkiSync.connection.gatewayId || undefined,
projectId: pkiSync.connection.projectId || undefined,
isPlatformManagedCredentials: pkiSync.connection.isPlatformManagedCredentials || undefined
}
};
};
const listPkiSyncsByProjectId = async ({ projectId }: TListPkiSyncsByProjectId, actor: OrgServiceActor) => {
@@ -245,75 +302,78 @@ export const pkiSyncServiceFactory = ({
projectId
});
ForbiddenError.from(permission).throwUnlessCan(
ProjectPermissionPkiSyncActions.Read,
subject(ProjectPermissionSub.PkiSyncs, { projectId })
);
ForbiddenError.from(permission).throwUnlessCan(ProjectPermissionPkiSyncActions.Read, ProjectPermissionSub.PkiSyncs);
const pkiSyncs = await pkiSyncDAL.findByProjectId(projectId);
return pkiSyncs;
};
const listPkiSyncsBySubscriberId = async ({ subscriberId }: TListPkiSyncsBySubscriberId) => {
const pkiSyncs = await pkiSyncDAL.findBySubscriberId(subscriberId);
return pkiSyncs;
return pkiSyncs as TPkiSync[];
};
const findPkiSyncById = async ({ id, projectId }: TFindPkiSyncByIdDTO, actor: OrgServiceActor) => {
const { permission } = await permissionService.getProjectPermission({
actor: actor.type,
actorId: actor.id,
actorAuthMethod: actor.authMethod,
actorOrgId: actor.orgId,
actionProjectType: ActionProjectType.CertificateManager,
projectId
});
const pkiSync = await pkiSyncDAL.findByIdAndProjectId(id, projectId);
const pkiSync = await pkiSyncDAL.findById(id);
if (!pkiSync)
throw new NotFoundError({
message: `Could not find PKI Sync with ID "${id}"`
});
ForbiddenError.from(permission).throwUnlessCan(
ProjectPermissionPkiSyncActions.Read,
subject(ProjectPermissionSub.PkiSyncs, {
projectId,
subscriberId: pkiSync.subscriberId
})
);
if (projectId && pkiSync.projectId !== projectId) {
throw new NotFoundError({
message: `Could not find PKI Sync with ID "${id}" in project "${projectId}"`
});
}
return pkiSync;
};
const findPkiSyncByName = async ({ name, projectId }: TFindPkiSyncByNameDTO) => {
const pkiSync = await pkiSyncDAL.findByNameAndProjectId(name, projectId);
if (!pkiSync) throw new NotFoundError({ message: "PKI sync not found" });
return pkiSync;
};
const triggerPkiSyncSyncCertificatesById = async (
{ id, projectId }: Omit<TTriggerPkiSyncSyncCertificatesByIdDTO, "auditLogInfo">,
actor: OrgServiceActor
) => {
const { permission } = await permissionService.getProjectPermission({
actor: actor.type,
actorId: actor.id,
actorAuthMethod: actor.authMethod,
actorOrgId: actor.orgId,
actionProjectType: ActionProjectType.CertificateManager,
projectId
projectId: pkiSync.projectId
});
const pkiSync = await pkiSyncDAL.findByIdAndProjectId(id, projectId);
let findSubscriber;
if (pkiSync.subscriberId) {
findSubscriber = await pkiSubscriberDAL.findById(pkiSync.subscriberId);
}
ForbiddenError.from(permission).throwUnlessCan(
ProjectPermissionPkiSyncActions.Read,
findSubscriber
? subject(ProjectPermissionSub.PkiSyncs, { subscriberName: findSubscriber.name })
: ProjectPermissionSub.PkiSyncs
);
return pkiSync as TPkiSync;
};
const triggerPkiSyncSyncCertificatesById = async (
{ id }: Omit<TTriggerPkiSyncSyncCertificatesByIdDTO, "auditLogInfo" | "projectId">,
actor: OrgServiceActor
) => {
const existingSync = await pkiSyncDAL.findById(id);
if (!existingSync) throw new NotFoundError({ message: "PKI sync not found" });
const { permission } = await permissionService.getProjectPermission({
actor: actor.type,
actorId: actor.id,
actorAuthMethod: actor.authMethod,
actorOrgId: actor.orgId,
actionProjectType: ActionProjectType.CertificateManager,
projectId: existingSync.projectId
});
const pkiSync = await pkiSyncDAL.findByIdAndProjectId(id, existingSync.projectId);
if (!pkiSync) throw new NotFoundError({ message: "PKI sync not found" });
let syncSubscriber;
if (pkiSync.subscriberId) {
syncSubscriber = await pkiSubscriberDAL.findById(pkiSync.subscriberId);
}
ForbiddenError.from(permission).throwUnlessCan(
ProjectPermissionPkiSyncActions.SyncCertificates,
subject(ProjectPermissionSub.PkiSyncs, {
projectId,
subscriberId: pkiSync.subscriberId
})
syncSubscriber
? subject(ProjectPermissionSub.PkiSyncs, { subscriberName: syncSubscriber.name })
: ProjectPermissionSub.PkiSyncs
);
await pkiSyncQueue.queuePkiSyncSyncCertificatesById({ syncId: id });
@@ -322,27 +382,42 @@ export const pkiSyncServiceFactory = ({
};
const triggerPkiSyncImportCertificatesById = async (
{ id, projectId }: Omit<TTriggerPkiSyncImportCertificatesByIdDTO, "auditLogInfo">,
{ id }: Omit<TTriggerPkiSyncImportCertificatesByIdDTO, "auditLogInfo" | "projectId">,
actor: OrgServiceActor
) => {
const existingSync = await pkiSyncDAL.findById(id);
if (!existingSync) throw new NotFoundError({ message: "PKI sync not found" });
const { permission } = await permissionService.getProjectPermission({
actor: actor.type,
actorId: actor.id,
actorAuthMethod: actor.authMethod,
actorOrgId: actor.orgId,
actionProjectType: ActionProjectType.CertificateManager,
projectId
projectId: existingSync.projectId
});
const pkiSync = await pkiSyncDAL.findByIdAndProjectId(id, projectId);
const pkiSync = await pkiSyncDAL.findByIdAndProjectId(id, existingSync.projectId);
if (!pkiSync) throw new NotFoundError({ message: "PKI sync not found" });
// Check if the PKI sync destination supports importing certificates
const syncOptions = listPkiSyncOptions().find((option) => option.destination === pkiSync.destination);
if (!syncOptions?.canImportCertificates) {
throw new BadRequestError({
message: `Certificate import is not supported for ${pkiSync.destination} PKI sync destination`
});
}
let importSubscriber;
if (pkiSync.subscriberId) {
importSubscriber = await pkiSubscriberDAL.findById(pkiSync.subscriberId);
}
ForbiddenError.from(permission).throwUnlessCan(
ProjectPermissionPkiSyncActions.ImportCertificates,
subject(ProjectPermissionSub.PkiSyncs, {
projectId,
subscriberId: pkiSync.subscriberId
})
importSubscriber
? subject(ProjectPermissionSub.PkiSyncs, { subscriberName: importSubscriber.name })
: ProjectPermissionSub.PkiSyncs
);
await pkiSyncQueue.queuePkiSyncImportCertificatesById({ syncId: id });
@@ -351,27 +426,34 @@ export const pkiSyncServiceFactory = ({
};
const triggerPkiSyncRemoveCertificatesById = async (
{ id, projectId }: Omit<TTriggerPkiSyncRemoveCertificatesByIdDTO, "auditLogInfo">,
{ id }: Omit<TTriggerPkiSyncRemoveCertificatesByIdDTO, "auditLogInfo" | "projectId">,
actor: OrgServiceActor
) => {
const existingSync = await pkiSyncDAL.findById(id);
if (!existingSync) throw new NotFoundError({ message: "PKI sync not found" });
const { permission } = await permissionService.getProjectPermission({
actor: actor.type,
actorId: actor.id,
actorAuthMethod: actor.authMethod,
actorOrgId: actor.orgId,
actionProjectType: ActionProjectType.CertificateManager,
projectId
projectId: existingSync.projectId
});
const pkiSync = await pkiSyncDAL.findByIdAndProjectId(id, projectId);
const pkiSync = await pkiSyncDAL.findByIdAndProjectId(id, existingSync.projectId);
if (!pkiSync) throw new NotFoundError({ message: "PKI sync not found" });
let removeSubscriber;
if (pkiSync.subscriberId) {
removeSubscriber = await pkiSubscriberDAL.findById(pkiSync.subscriberId);
}
ForbiddenError.from(permission).throwUnlessCan(
ProjectPermissionPkiSyncActions.RemoveCertificates,
subject(ProjectPermissionSub.PkiSyncs, {
projectId,
subscriberId: pkiSync.subscriberId
})
removeSubscriber
? subject(ProjectPermissionSub.PkiSyncs, { subscriberName: removeSubscriber.name })
: ProjectPermissionSub.PkiSyncs
);
await pkiSyncQueue.queuePkiSyncRemoveCertificatesById({ syncId: id });
@@ -388,9 +470,7 @@ export const pkiSyncServiceFactory = ({
updatePkiSync,
deletePkiSync,
listPkiSyncsByProjectId,
listPkiSyncsBySubscriberId,
findPkiSyncById,
findPkiSyncByName,
triggerPkiSyncSyncCertificatesById,
triggerPkiSyncImportCertificatesById,
triggerPkiSyncRemoveCertificatesById,

View File

@@ -13,7 +13,6 @@ export type TPkiSync = {
description?: string;
destination: PkiSync;
isAutoSyncEnabled: boolean;
version: number;
destinationConfig: Record<string, unknown>;
syncOptions: Record<string, unknown>;
projectId: string;
@@ -33,11 +32,23 @@ export type TPkiSync = {
lastRemoveJobId?: string;
lastRemoveMessage?: string;
lastRemovedAt?: Date;
};
export type TPkiSyncListItem = TPkiSync & {
appConnectionName: string;
appConnectionApp: string;
connection: {
id: string;
name: string;
app: string;
encryptedCredentials: unknown;
orgId: string;
projectId?: string;
method: string;
description?: string;
version: number;
gatewayId?: string;
createdAt: Date;
updatedAt: Date;
isPlatformManagedCredentials?: boolean;
};
};
export type TPkiSyncWithCredentials = TPkiSync & {
@@ -50,6 +61,11 @@ export type TPkiSyncWithCredentials = TPkiSync & {
};
};
export type TPkiSyncListItem = TPkiSync & {
appConnectionName: string;
appConnectionApp: string;
};
export type TCertificateMap = Record<string, { cert: string; privateKey: string }>;
export type TCreatePkiSyncDTO = {
@@ -68,7 +84,7 @@ export type TCreatePkiSyncDTO = {
export type TUpdatePkiSyncDTO = {
id: string;
projectId: string;
projectId?: string;
name?: string;
description?: string;
isAutoSyncEnabled?: boolean;
@@ -82,7 +98,7 @@ export type TUpdatePkiSyncDTO = {
export type TDeletePkiSyncDTO = {
id: string;
projectId: string;
projectId?: string;
auditLogInfo: AuditLogInfo;
};
@@ -90,51 +106,29 @@ export type TListPkiSyncsByProjectId = {
projectId: string;
};
export type TListPkiSyncsBySubscriberId = {
subscriberId: string;
};
export type TFindPkiSyncByIdDTO = {
id: string;
projectId: string;
};
export type TFindPkiSyncByNameDTO = {
name: string;
projectId: string;
projectId?: string;
};
export type TTriggerPkiSyncSyncCertificatesByIdDTO = {
id: string;
projectId: string;
projectId?: string;
auditLogInfo: AuditLogInfo;
};
export type TTriggerPkiSyncImportCertificatesByIdDTO = {
id: string;
projectId: string;
projectId?: string;
auditLogInfo: AuditLogInfo;
};
export type TTriggerPkiSyncRemoveCertificatesByIdDTO = {
id: string;
projectId: string;
projectId?: string;
auditLogInfo: AuditLogInfo;
};
export enum PkiSyncStatus {
Pending = "pending",
Running = "running",
Succeeded = "succeeded",
Failed = "failed"
}
export enum PkiSyncAction {
SyncCertificates = "sync-certificates",
ImportCertificates = "import-certificates",
RemoveCertificates = "remove-certificates"
}
export type TPkiSyncRaw = NonNullable<Awaited<ReturnType<TPkiSyncDALFactory["findById"]>>>;
export type TQueuePkiSyncSyncCertificatesByIdDTO = {
@@ -154,12 +148,6 @@ export type TQueuePkiSyncRemoveCertificatesByIdDTO = {
deleteSyncOnComplete?: boolean;
};
export type TQueueSendPkiSyncActionFailedNotificationsDTO = {
pkiSync: TPkiSyncRaw;
auditLogInfo?: AuditLogInfo;
action: PkiSyncAction;
};
export type TPkiSyncSyncCertificatesDTO = Job<
TQueuePkiSyncSyncCertificatesByIdDTO,
void,
@@ -175,9 +163,3 @@ export type TPkiSyncRemoveCertificatesDTO = Job<
void,
QueueJobs.PkiSyncRemoveCertificates
>;
export type TSendPkiSyncFailedNotificationsJobDTO = Job<
TQueueSendPkiSyncActionFailedNotificationsDTO,
void,
QueueJobs.PkiSyncSendActionFailedNotifications
>;