Merge remote-tracking branch 'origin/main' into misc/add-endpoint-for-fetching-latest-active-bundle-from-profile

This commit is contained in:
Sheen Capadngan
2025-11-10 21:02:29 +08:00
137 changed files with 6777 additions and 522 deletions
+2
View File
@@ -101,6 +101,7 @@ import { TOfflineUsageReportServiceFactory } from "@app/services/offline-usage-r
import { TOrgServiceFactory } from "@app/services/org/org-service";
import { TOrgAdminServiceFactory } from "@app/services/org-admin/org-admin-service";
import { TPkiAlertServiceFactory } from "@app/services/pki-alert/pki-alert-service";
import { TPkiAlertV2ServiceFactory } from "@app/services/pki-alert-v2/pki-alert-v2-service";
import { TPkiCollectionServiceFactory } from "@app/services/pki-collection/pki-collection-service";
import { TPkiSubscriberServiceFactory } from "@app/services/pki-subscriber/pki-subscriber-service";
import { TPkiSyncServiceFactory } from "@app/services/pki-sync/pki-sync-service";
@@ -355,6 +356,7 @@ declare module "fastify" {
role: TRoleServiceFactory;
convertor: TConvertorServiceFactory;
subOrganization: TSubOrgServiceFactory;
pkiAlertV2: TPkiAlertV2ServiceFactory;
};
// this is exclusive use for middlewares in which we need to inject data
// everywhere else access using service layer
+28
View File
@@ -284,9 +284,21 @@ import {
TPkiAcmeOrders,
TPkiAcmeOrdersInsert,
TPkiAcmeOrdersUpdate,
TPkiAlertChannels,
TPkiAlertChannelsInsert,
TPkiAlertChannelsUpdate,
TPkiAlertHistory,
TPkiAlertHistoryCertificate,
TPkiAlertHistoryCertificateInsert,
TPkiAlertHistoryCertificateUpdate,
TPkiAlertHistoryInsert,
TPkiAlertHistoryUpdate,
TPkiAlerts,
TPkiAlertsInsert,
TPkiAlertsUpdate,
TPkiAlertsV2,
TPkiAlertsV2Insert,
TPkiAlertsV2Update,
TPkiApiEnrollmentConfigs,
TPkiApiEnrollmentConfigsInsert,
TPkiApiEnrollmentConfigsUpdate,
@@ -769,6 +781,22 @@ declare module "knex/types/tables" {
TCertificateSecretsUpdate
>;
[TableName.PkiAlert]: KnexOriginal.CompositeTableType<TPkiAlerts, TPkiAlertsInsert, TPkiAlertsUpdate>;
[TableName.PkiAlertsV2]: KnexOriginal.CompositeTableType<TPkiAlertsV2, TPkiAlertsV2Insert, TPkiAlertsV2Update>;
[TableName.PkiAlertChannels]: KnexOriginal.CompositeTableType<
TPkiAlertChannels,
TPkiAlertChannelsInsert,
TPkiAlertChannelsUpdate
>;
[TableName.PkiAlertHistory]: KnexOriginal.CompositeTableType<
TPkiAlertHistory,
TPkiAlertHistoryInsert,
TPkiAlertHistoryUpdate
>;
[TableName.PkiAlertHistoryCertificate]: KnexOriginal.CompositeTableType<
TPkiAlertHistoryCertificate,
TPkiAlertHistoryCertificateInsert,
TPkiAlertHistoryCertificateUpdate
>;
[TableName.PkiCollection]: KnexOriginal.CompositeTableType<
TPkiCollections,
TPkiCollectionsInsert,
@@ -0,0 +1,87 @@
import { Knex } from "knex";
import { TableName } from "../schemas";
export async function up(knex: Knex): Promise<void> {
if (!(await knex.schema.hasTable(TableName.PkiAlertsV2))) {
await knex.schema.createTable(TableName.PkiAlertsV2, (t) => {
t.uuid("id", { primaryKey: true }).defaultTo(knex.fn.uuid());
t.string("name").notNullable();
t.text("description").nullable();
t.string("eventType").notNullable();
t.string("alertBefore").nullable();
t.jsonb("filters").nullable();
t.boolean("enabled").defaultTo(true);
t.string("projectId").notNullable();
t.timestamps(true, true, true);
t.foreign("projectId").references("id").inTable(TableName.Project).onDelete("CASCADE");
t.index("projectId");
t.unique(["name", "projectId"]);
});
}
if (!(await knex.schema.hasTable(TableName.PkiAlertChannels))) {
await knex.schema.createTable(TableName.PkiAlertChannels, (t) => {
t.uuid("id", { primaryKey: true }).defaultTo(knex.fn.uuid());
t.uuid("alertId").notNullable();
t.string("channelType").notNullable();
t.jsonb("config").notNullable();
t.boolean("enabled").defaultTo(true);
t.timestamps(true, true, true);
t.foreign("alertId").references("id").inTable(TableName.PkiAlertsV2).onDelete("CASCADE");
t.index("alertId");
t.index("channelType");
t.index("enabled");
});
}
if (!(await knex.schema.hasTable(TableName.PkiAlertHistory))) {
await knex.schema.createTable(TableName.PkiAlertHistory, (t) => {
t.uuid("id", { primaryKey: true }).defaultTo(knex.fn.uuid());
t.uuid("alertId").notNullable();
t.timestamp("triggeredAt").defaultTo(knex.fn.now());
t.boolean("hasNotificationSent").defaultTo(false);
t.text("notificationError").nullable();
t.timestamps(true, true, true);
t.foreign("alertId").references("id").inTable(TableName.PkiAlertsV2).onDelete("CASCADE");
t.index("alertId");
t.index("triggeredAt");
});
}
if (!(await knex.schema.hasTable(TableName.PkiAlertHistoryCertificate))) {
await knex.schema.createTable(TableName.PkiAlertHistoryCertificate, (t) => {
t.uuid("id", { primaryKey: true }).defaultTo(knex.fn.uuid());
t.uuid("alertHistoryId").notNullable();
t.uuid("certificateId").notNullable();
t.timestamps(true, true, true);
t.foreign("alertHistoryId").references("id").inTable(TableName.PkiAlertHistory).onDelete("CASCADE");
t.foreign("certificateId").references("id").inTable(TableName.Certificate).onDelete("CASCADE");
t.index("alertHistoryId");
t.index("certificateId");
t.unique(["alertHistoryId", "certificateId"]);
});
}
}
export async function down(knex: Knex): Promise<void> {
if (await knex.schema.hasTable(TableName.PkiAlertHistoryCertificate)) {
await knex.schema.dropTable(TableName.PkiAlertHistoryCertificate);
}
if (await knex.schema.hasTable(TableName.PkiAlertHistory)) {
await knex.schema.dropTable(TableName.PkiAlertHistory);
}
if (await knex.schema.hasTable(TableName.PkiAlertChannels)) {
await knex.schema.dropTable(TableName.PkiAlertChannels);
}
if (await knex.schema.hasTable(TableName.PkiAlertsV2)) {
await knex.schema.dropTable(TableName.PkiAlertsV2);
}
}
+4
View File
@@ -98,7 +98,11 @@ export * from "./pki-acme-challenges";
export * from "./pki-acme-enrollment-configs";
export * from "./pki-acme-order-auths";
export * from "./pki-acme-orders";
export * from "./pki-alert-channels";
export * from "./pki-alert-history";
export * from "./pki-alert-history-certificate";
export * from "./pki-alerts";
export * from "./pki-alerts-v2";
export * from "./pki-api-enrollment-configs";
export * from "./pki-certificate-profiles";
export * from "./pki-certificate-templates-v2";
+4
View File
@@ -30,6 +30,10 @@ export enum TableName {
PkiAcmeEnrollmentConfig = "pki_acme_enrollment_configs",
PkiSubscriber = "pki_subscribers",
PkiAlert = "pki_alerts",
PkiAlertsV2 = "pki_alerts_v2",
PkiAlertChannels = "pki_alert_channels",
PkiAlertHistory = "pki_alert_history",
PkiAlertHistoryCertificate = "pki_alert_history_certificate",
PkiCollection = "pki_collections",
PkiCollectionItem = "pki_collection_items",
Groups = "groups",
@@ -0,0 +1,22 @@
// Code generated by automation script, DO NOT EDIT.
// Automated by pulling database and generating zod schema
// To update. Just run npm run generate:schema
// Written by akhilmhdh.
import { z } from "zod";
import { TImmutableDBKeys } from "./models";
export const PkiAlertChannelsSchema = z.object({
id: z.string().uuid(),
alertId: z.string().uuid(),
channelType: z.string(),
config: z.unknown(),
enabled: z.boolean().default(true).nullable().optional(),
createdAt: z.date(),
updatedAt: z.date()
});
export type TPkiAlertChannels = z.infer<typeof PkiAlertChannelsSchema>;
export type TPkiAlertChannelsInsert = Omit<z.input<typeof PkiAlertChannelsSchema>, TImmutableDBKeys>;
export type TPkiAlertChannelsUpdate = Partial<Omit<z.input<typeof PkiAlertChannelsSchema>, TImmutableDBKeys>>;
@@ -0,0 +1,25 @@
// Code generated by automation script, DO NOT EDIT.
// Automated by pulling database and generating zod schema
// To update. Just run npm run generate:schema
// Written by akhilmhdh.
import { z } from "zod";
import { TImmutableDBKeys } from "./models";
export const PkiAlertHistoryCertificateSchema = z.object({
id: z.string().uuid(),
alertHistoryId: z.string().uuid(),
certificateId: z.string().uuid(),
createdAt: z.date(),
updatedAt: z.date()
});
export type TPkiAlertHistoryCertificate = z.infer<typeof PkiAlertHistoryCertificateSchema>;
export type TPkiAlertHistoryCertificateInsert = Omit<
z.input<typeof PkiAlertHistoryCertificateSchema>,
TImmutableDBKeys
>;
export type TPkiAlertHistoryCertificateUpdate = Partial<
Omit<z.input<typeof PkiAlertHistoryCertificateSchema>, TImmutableDBKeys>
>;
@@ -0,0 +1,22 @@
// Code generated by automation script, DO NOT EDIT.
// Automated by pulling database and generating zod schema
// To update. Just run npm run generate:schema
// Written by akhilmhdh.
import { z } from "zod";
import { TImmutableDBKeys } from "./models";
export const PkiAlertHistorySchema = z.object({
id: z.string().uuid(),
alertId: z.string().uuid(),
triggeredAt: z.date().nullable().optional(),
hasNotificationSent: z.boolean().default(false).nullable().optional(),
notificationError: z.string().nullable().optional(),
createdAt: z.date(),
updatedAt: z.date()
});
export type TPkiAlertHistory = z.infer<typeof PkiAlertHistorySchema>;
export type TPkiAlertHistoryInsert = Omit<z.input<typeof PkiAlertHistorySchema>, TImmutableDBKeys>;
export type TPkiAlertHistoryUpdate = Partial<Omit<z.input<typeof PkiAlertHistorySchema>, TImmutableDBKeys>>;
+25
View File
@@ -0,0 +1,25 @@
// Code generated by automation script, DO NOT EDIT.
// Automated by pulling database and generating zod schema
// To update. Just run npm run generate:schema
// Written by akhilmhdh.
import { z } from "zod";
import { TImmutableDBKeys } from "./models";
export const PkiAlertsV2Schema = z.object({
id: z.string().uuid(),
name: z.string(),
description: z.string().nullable().optional(),
eventType: z.string(),
alertBefore: z.string().nullable().optional(),
filters: z.unknown().nullable().optional(),
enabled: z.boolean().default(true).nullable().optional(),
projectId: z.string(),
createdAt: z.date(),
updatedAt: z.date()
});
export type TPkiAlertsV2 = z.infer<typeof PkiAlertsV2Schema>;
export type TPkiAlertsV2Insert = Omit<z.input<typeof PkiAlertsV2Schema>, TImmutableDBKeys>;
export type TPkiAlertsV2Update = Partial<Omit<z.input<typeof PkiAlertsV2Schema>, TImmutableDBKeys>>;
@@ -36,6 +36,7 @@ import { CertExtendedKeyUsage, CertKeyAlgorithm, CertKeyUsage } from "@app/servi
import { CaStatus } from "@app/services/certificate-authority/certificate-authority-enums";
import { TIdentityTrustedIp } from "@app/services/identity/identity-types";
import { TAllowedFields } from "@app/services/identity-ldap-auth/identity-ldap-auth-types";
import { PkiAlertEventType } from "@app/services/pki-alert-v2/pki-alert-v2-types";
import { PkiItemType } from "@app/services/pki-collection/pki-collection-types";
import { SecretSync, SecretSyncImportBehavior } from "@app/services/secret-sync/secret-sync-enums";
import {
@@ -2318,10 +2319,11 @@ interface CreatePkiAlert {
type: EventType.CREATE_PKI_ALERT;
metadata: {
pkiAlertId: string;
pkiCollectionId: string;
pkiCollectionId?: string;
name: string;
alertBeforeDays: number;
recipientEmails: string;
alertBefore: string;
eventType: PkiAlertEventType;
recipientEmails?: string;
};
}
interface GetPkiAlert {
@@ -2337,7 +2339,8 @@ interface UpdatePkiAlert {
pkiAlertId: string;
pkiCollectionId?: string;
name?: string;
alertBeforeDays?: number;
alertBefore?: string;
eventType?: PkiAlertEventType;
recipientEmails?: string;
};
}
+4 -1
View File
@@ -2309,7 +2309,10 @@ export const AppConnections = {
code: "The OAuth code to use to connect with Azure Client Secrets.",
tenantId: "The Tenant ID to use to connect with Azure Client Secrets.",
clientId: "The Client ID to use to connect with Azure Client Secrets.",
clientSecret: "The Client Secret to use to connect with Azure Client Secrets."
clientSecret: "The Client Secret to use to connect with Azure Client Secrets.",
certificateBody: "The certificate body in PEM format to use to connect with Azure Client Secrets.",
privateKey:
"The private key to use to connect with Azure Client Secrets. This is never transmitted to Azure and is only used to sign the Azure client assertion with."
},
AZURE_DEVOPS: {
code: "The OAuth code to use to connect with Azure DevOps.",
+6
View File
@@ -51,6 +51,7 @@ export enum QueueName {
AuditLogPrune = "audit-log-prune",
DailyResourceCleanUp = "daily-resource-cleanup",
DailyExpiringPkiItemAlert = "daily-expiring-pki-item-alert",
DailyPkiAlertV2Processing = "daily-pki-alert-v2-processing",
PkiSyncCleanup = "pki-sync-cleanup",
PkiSubscriber = "pki-subscriber",
TelemetryInstanceStats = "telemtry-self-hosted-stats",
@@ -90,6 +91,7 @@ export enum QueueJobs {
AuditLogPrune = "audit-log-prune-job",
DailyResourceCleanUp = "daily-resource-cleanup-job",
DailyExpiringPkiItemAlert = "daily-expiring-pki-item-alert",
DailyPkiAlertV2Processing = "daily-pki-alert-v2-processing",
PkiSyncCleanup = "pki-sync-cleanup-job",
SecWebhook = "secret-webhook-trigger",
TelemetryInstanceStats = "telemetry-self-hosted-stats",
@@ -159,6 +161,10 @@ export type TQueueJobTypes = {
name: QueueJobs.DailyExpiringPkiItemAlert;
payload: undefined;
};
[QueueName.DailyPkiAlertV2Processing]: {
name: QueueJobs.DailyPkiAlertV2Processing;
payload: undefined;
};
[QueueName.PkiSyncCleanup]: {
name: QueueJobs.PkiSyncCleanup;
payload: undefined;
+26 -1
View File
@@ -275,6 +275,11 @@ import { pamAccountRotationServiceFactory } from "@app/services/pam-account-rota
import { dailyExpiringPkiItemAlertQueueServiceFactory } from "@app/services/pki-alert/expiring-pki-item-alert-queue";
import { pkiAlertDALFactory } from "@app/services/pki-alert/pki-alert-dal";
import { pkiAlertServiceFactory } from "@app/services/pki-alert/pki-alert-service";
import { pkiAlertChannelDALFactory } from "@app/services/pki-alert-v2/pki-alert-channel-dal";
import { pkiAlertHistoryDALFactory } from "@app/services/pki-alert-v2/pki-alert-history-dal";
import { pkiAlertV2DALFactory } from "@app/services/pki-alert-v2/pki-alert-v2-dal";
import { pkiAlertV2QueueServiceFactory } from "@app/services/pki-alert-v2/pki-alert-v2-queue";
import { pkiAlertV2ServiceFactory } from "@app/services/pki-alert-v2/pki-alert-v2-service";
import { pkiCollectionDALFactory } from "@app/services/pki-collection/pki-collection-dal";
import { pkiCollectionItemDALFactory } from "@app/services/pki-collection/pki-collection-item-dal";
import { pkiCollectionServiceFactory } from "@app/services/pki-collection/pki-collection-service";
@@ -555,6 +560,9 @@ export const registerRoutes = async (
const additionalPrivilegeDAL = additionalPrivilegeDALFactory(db);
const membershipRoleDAL = membershipRoleDALFactory(db);
const roleDAL = roleDALFactory(db);
const pkiAlertHistoryDAL = pkiAlertHistoryDALFactory(db);
const pkiAlertChannelDAL = pkiAlertChannelDALFactory(db);
const pkiAlertV2DAL = pkiAlertV2DALFactory(db);
const vaultExternalMigrationConfigDAL = vaultExternalMigrationConfigDALFactory(db);
@@ -1822,6 +1830,21 @@ export const registerRoutes = async (
groupDAL
});
const pkiAlertV2Service = pkiAlertV2ServiceFactory({
pkiAlertV2DAL,
pkiAlertChannelDAL,
pkiAlertHistoryDAL,
permissionService,
smtpService
});
const pkiAlertV2Queue = pkiAlertV2QueueServiceFactory({
queueService,
pkiAlertV2Service,
pkiAlertV2DAL,
pkiAlertHistoryDAL
});
const dynamicSecretProviders = buildDynamicSecretProviders({
gatewayService,
gatewayV2Service
@@ -2395,6 +2418,7 @@ export const registerRoutes = async (
await dailyReminderQueueService.startSecretReminderMigrationJob();
await dailyExpiringPkiItemAlert.startSendingAlerts();
await pkiSubscriberQueue.startDailyAutoRenewalJob();
await pkiAlertV2Queue.init();
await certificateV3Queue.init();
await kmsService.startService(hsmStatus);
await microsoftTeamsService.start();
@@ -2527,7 +2551,8 @@ export const registerRoutes = async (
role: roleService,
additionalPrivilege: additionalPrivilegeService,
identityProject: identityProjectService,
convertor: convertorService
convertor: convertorService,
pkiAlertV2: pkiAlertV2Service
});
const cronJobs: CronJob[] = [];
@@ -6,6 +6,7 @@ import { ALERTS, ApiDocsTags } from "@app/lib/api-docs";
import { readLimit, writeLimit } from "@app/server/config/rateLimiter";
import { verifyAuth } from "@app/server/plugins/auth/verify-auth";
import { AuthMode } from "@app/services/auth/auth-type";
import { PkiAlertEventType } from "@app/services/pki-alert-v2/pki-alert-v2-types";
export const registerPkiAlertRouter = async (server: FastifyZodProvider) => {
server.route({
@@ -52,7 +53,8 @@ export const registerPkiAlertRouter = async (server: FastifyZodProvider) => {
pkiAlertId: alert.id,
pkiCollectionId: alert.pkiCollectionId,
name: alert.name,
alertBeforeDays: alert.alertBeforeDays,
alertBefore: alert.alertBeforeDays.toString(),
eventType: PkiAlertEventType.EXPIRATION,
recipientEmails: alert.recipientEmails
}
}
@@ -152,7 +154,8 @@ export const registerPkiAlertRouter = async (server: FastifyZodProvider) => {
pkiAlertId: alert.id,
pkiCollectionId: alert.pkiCollectionId,
name: alert.name,
alertBeforeDays: alert.alertBeforeDays,
alertBefore: alert.alertBeforeDays.toString(),
eventType: PkiAlertEventType.EXPIRATION,
recipientEmails: alert.recipientEmails
}
}
+2
View File
@@ -8,6 +8,7 @@ import { registerIdentityOrgRouter } from "./identity-org-router";
import { registerMfaRouter } from "./mfa-router";
import { registerOrgRouter } from "./organization-router";
import { registerPasswordRouter } from "./password-router";
import { registerPkiAlertRouter } from "./pki-alert-router";
import { registerPkiTemplatesRouter } from "./pki-templates-router";
import { registerSecretFolderRouter } from "./secret-folder-router";
import { registerSecretImportRouter } from "./secret-import-router";
@@ -26,6 +27,7 @@ export const registerV2Routes = async (server: FastifyZodProvider) => {
async (pkiRouter) => {
await pkiRouter.register(registerCaRouter, { prefix: "/ca" });
await pkiRouter.register(registerPkiTemplatesRouter, { prefix: "/certificate-templates" });
await pkiRouter.register(registerPkiAlertRouter, { prefix: "/alerts" });
},
{ prefix: "/pki" }
);
@@ -0,0 +1,443 @@
import { z } from "zod";
import { EventType } from "@app/ee/services/audit-log/audit-log-types";
import { ApiDocsTags } from "@app/lib/api-docs";
import { readLimit, writeLimit } from "@app/server/config/rateLimiter";
import { verifyAuth } from "@app/server/plugins/auth/verify-auth";
import { AuthMode } from "@app/services/auth/auth-type";
import {
CreatePkiAlertV2Schema,
createSecureAlertBeforeValidator,
PkiAlertChannelType,
PkiAlertEventType,
PkiFilterRuleSchema,
UpdatePkiAlertV2Schema
} from "@app/services/pki-alert-v2/pki-alert-v2-types";
export const registerPkiAlertRouter = async (server: FastifyZodProvider) => {
server.route({
method: "POST",
url: "/",
config: {
rateLimit: writeLimit
},
onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]),
schema: {
description: "Create a new PKI alert",
tags: [ApiDocsTags.PkiAlerting],
body: CreatePkiAlertV2Schema.extend({
projectId: z.string().uuid().describe("Project ID")
}),
response: {
200: z.object({
alert: z.object({
id: z.string().uuid(),
name: z.string(),
description: z.string().nullable(),
eventType: z.nativeEnum(PkiAlertEventType),
alertBefore: z.string(),
filters: z.array(PkiFilterRuleSchema),
enabled: z.boolean(),
projectId: z.string().uuid(),
channels: z.array(
z.object({
id: z.string().uuid(),
channelType: z.nativeEnum(PkiAlertChannelType),
config: z.record(z.any()),
enabled: z.boolean(),
createdAt: z.date(),
updatedAt: z.date()
})
),
createdAt: z.date(),
updatedAt: z.date()
})
})
}
},
handler: async (req) => {
const alert = await server.services.pkiAlertV2.createAlert({
actor: req.permission.type,
actorId: req.permission.id,
actorAuthMethod: req.permission.authMethod,
actorOrgId: req.permission.orgId,
...req.body
});
await server.services.auditLog.createAuditLog({
...req.auditLogInfo,
projectId: req.body.projectId,
event: {
type: EventType.CREATE_PKI_ALERT,
metadata: {
pkiAlertId: alert.id,
name: alert.name,
eventType: alert.eventType,
alertBefore: alert.alertBefore
}
}
});
return { alert };
}
});
server.route({
method: "GET",
url: "/",
config: {
rateLimit: readLimit
},
onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]),
schema: {
description: "List PKI alerts for a project",
tags: [ApiDocsTags.PkiAlerting],
querystring: z.object({
projectId: z.string().uuid(),
search: z.string().optional(),
eventType: z.nativeEnum(PkiAlertEventType).optional(),
enabled: z.coerce.boolean().optional(),
limit: z.coerce.number().min(1).max(100).default(20),
offset: z.coerce.number().min(0).default(0)
}),
response: {
200: z.object({
alerts: z.array(
z.object({
id: z.string().uuid(),
name: z.string(),
description: z.string().nullable(),
eventType: z.nativeEnum(PkiAlertEventType),
alertBefore: z.string(),
filters: z.array(PkiFilterRuleSchema),
enabled: z.boolean(),
channels: z.array(
z.object({
id: z.string().uuid(),
channelType: z.nativeEnum(PkiAlertChannelType),
config: z.record(z.any()),
enabled: z.boolean(),
createdAt: z.date(),
updatedAt: z.date()
})
),
createdAt: z.date(),
updatedAt: z.date()
})
),
total: z.number()
})
}
},
handler: async (req) => {
const alerts = await server.services.pkiAlertV2.listAlerts({
actor: req.permission.type,
actorId: req.permission.id,
actorAuthMethod: req.permission.authMethod,
actorOrgId: req.permission.orgId,
...req.query
});
return alerts;
}
});
server.route({
method: "GET",
url: "/:alertId",
config: {
rateLimit: readLimit
},
onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]),
schema: {
description: "Get a PKI alert by ID",
tags: [ApiDocsTags.PkiAlerting],
params: z.object({
alertId: z.string().uuid().describe("Alert ID")
}),
response: {
200: z.object({
alert: z.object({
id: z.string().uuid(),
name: z.string(),
description: z.string().nullable(),
eventType: z.nativeEnum(PkiAlertEventType),
alertBefore: z.string(),
filters: z.array(PkiFilterRuleSchema),
enabled: z.boolean(),
projectId: z.string().uuid(),
channels: z.array(
z.object({
id: z.string().uuid(),
channelType: z.nativeEnum(PkiAlertChannelType),
config: z.record(z.any()),
enabled: z.boolean(),
createdAt: z.date(),
updatedAt: z.date()
})
),
createdAt: z.date(),
updatedAt: z.date()
})
})
}
},
handler: async (req) => {
const alert = await server.services.pkiAlertV2.getAlertById({
alertId: req.params.alertId,
actor: req.permission.type,
actorId: req.permission.id,
actorAuthMethod: req.permission.authMethod,
actorOrgId: req.permission.orgId
});
await server.services.auditLog.createAuditLog({
...req.auditLogInfo,
projectId: alert.projectId,
event: {
type: EventType.GET_PKI_ALERT,
metadata: {
pkiAlertId: alert.id
}
}
});
return { alert };
}
});
server.route({
method: "PATCH",
url: "/:alertId",
config: {
rateLimit: writeLimit
},
onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]),
schema: {
description: "Update a PKI alert",
tags: [ApiDocsTags.PkiAlerting],
params: z.object({
alertId: z.string().uuid().describe("Alert ID")
}),
body: UpdatePkiAlertV2Schema,
response: {
200: z.object({
alert: z.object({
id: z.string().uuid(),
name: z.string(),
description: z.string().nullable(),
eventType: z.nativeEnum(PkiAlertEventType),
alertBefore: z.string(),
filters: z.array(PkiFilterRuleSchema),
enabled: z.boolean(),
projectId: z.string().uuid(),
channels: z.array(
z.object({
id: z.string().uuid(),
channelType: z.nativeEnum(PkiAlertChannelType),
config: z.record(z.any()),
enabled: z.boolean(),
createdAt: z.date(),
updatedAt: z.date()
})
),
createdAt: z.date(),
updatedAt: z.date()
})
})
}
},
handler: async (req) => {
const alert = await server.services.pkiAlertV2.updateAlert({
alertId: req.params.alertId,
actor: req.permission.type,
actorId: req.permission.id,
actorAuthMethod: req.permission.authMethod,
actorOrgId: req.permission.orgId,
...req.body
});
await server.services.auditLog.createAuditLog({
...req.auditLogInfo,
projectId: alert.projectId,
event: {
type: EventType.UPDATE_PKI_ALERT,
metadata: {
pkiAlertId: alert.id,
name: alert.name,
eventType: alert.eventType,
alertBefore: alert.alertBefore
}
}
});
return { alert };
}
});
server.route({
method: "DELETE",
url: "/:alertId",
config: {
rateLimit: writeLimit
},
onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]),
schema: {
description: "Delete a PKI alert",
tags: [ApiDocsTags.PkiAlerting],
params: z.object({
alertId: z.string().uuid().describe("Alert ID")
}),
response: {
200: z.object({
alert: z.object({
id: z.string().uuid(),
name: z.string(),
description: z.string().nullable(),
eventType: z.nativeEnum(PkiAlertEventType),
alertBefore: z.string(),
filters: z.array(PkiFilterRuleSchema),
enabled: z.boolean(),
projectId: z.string().uuid(),
channels: z.array(
z.object({
id: z.string().uuid(),
channelType: z.nativeEnum(PkiAlertChannelType),
config: z.record(z.any()),
enabled: z.boolean(),
createdAt: z.date(),
updatedAt: z.date()
})
),
createdAt: z.date(),
updatedAt: z.date()
})
})
}
},
handler: async (req) => {
const alert = await server.services.pkiAlertV2.deleteAlert({
alertId: req.params.alertId,
actor: req.permission.type,
actorId: req.permission.id,
actorAuthMethod: req.permission.authMethod,
actorOrgId: req.permission.orgId
});
await server.services.auditLog.createAuditLog({
...req.auditLogInfo,
projectId: alert.projectId,
event: {
type: EventType.DELETE_PKI_ALERT,
metadata: {
pkiAlertId: alert.id
}
}
});
return { alert };
}
});
server.route({
method: "GET",
url: "/:alertId/certificates",
config: {
rateLimit: readLimit
},
onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]),
schema: {
description: "List certificates that match an alert's filter rules",
tags: [ApiDocsTags.PkiAlerting],
params: z.object({
alertId: z.string().uuid().describe("Alert ID")
}),
querystring: z.object({
limit: z.coerce.number().min(1).max(100).default(20),
offset: z.coerce.number().min(0).default(0)
}),
response: {
200: z.object({
certificates: z.array(
z.object({
id: z.string().uuid(),
serialNumber: z.string(),
commonName: z.string(),
san: z.array(z.string()),
profileName: z.string().nullable(),
enrollmentType: z.string().nullable(),
notBefore: z.date(),
notAfter: z.date(),
status: z.string()
})
),
total: z.number()
})
}
},
handler: async (req) => {
const result = await server.services.pkiAlertV2.listMatchingCertificates({
alertId: req.params.alertId,
actor: req.permission.type,
actorId: req.permission.id,
actorAuthMethod: req.permission.authMethod,
actorOrgId: req.permission.orgId,
...req.query
});
return result;
}
});
server.route({
method: "POST",
url: "/preview/certificates",
config: {
rateLimit: writeLimit
},
onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]),
schema: {
description: "Preview certificates that would match the given filter rules",
tags: [ApiDocsTags.PkiAlerting],
body: z.object({
projectId: z.string().uuid().describe("Project ID"),
filters: z.array(PkiFilterRuleSchema),
alertBefore: z
.string()
.refine(createSecureAlertBeforeValidator(), "Must be in format like '30d', '1w', '3m', '1y'")
.describe("Alert timing (e.g., '30d', '1w')"),
limit: z.coerce.number().min(1).max(100).default(20),
offset: z.coerce.number().min(0).default(0)
}),
response: {
200: z.object({
certificates: z.array(
z.object({
id: z.string().uuid(),
serialNumber: z.string(),
commonName: z.string(),
san: z.array(z.string()),
profileName: z.string().nullable(),
enrollmentType: z.string().nullable(),
notBefore: z.date(),
notAfter: z.date(),
status: z.string()
})
),
total: z.number()
})
}
},
handler: async (req) => {
const result = await server.services.pkiAlertV2.listCurrentMatchingCertificates({
actor: req.permission.type,
actorId: req.permission.id,
actorAuthMethod: req.permission.authMethod,
actorOrgId: req.permission.orgId,
...req.body
});
return result;
}
});
};
@@ -1,4 +1,5 @@
export enum AzureClientSecretsConnectionMethod {
OAuth = "oauth",
ClientSecret = "client-secret"
ClientSecret = "client-secret",
Certificate = "certificate"
}
@@ -1,9 +1,14 @@
/* eslint-disable no-case-declarations */
import { AxiosError, AxiosResponse } from "axios";
import type { KeyObject } from "crypto";
import RE2 from "re2";
import { v4 as uuidv4 } from "uuid";
import { getConfig } from "@app/lib/config/env";
import { request } from "@app/lib/config/request";
import { crypto } from "@app/lib/crypto";
import { BadRequestError, InternalServerError, NotFoundError } from "@app/lib/errors";
import { logger } from "@app/lib/logger";
import {
decryptAppConnectionCredentials,
encryptAppConnectionCredentials,
@@ -17,11 +22,82 @@ import { AppConnection } from "../app-connection-enums";
import { AzureClientSecretsConnectionMethod } from "./azure-client-secrets-connection-enums";
import {
ExchangeCodeAzureResponse,
TAzureClientSecretsConnectionCertificateCredentials,
TAzureClientSecretsConnectionClientSecretCredentials,
TAzureClientSecretsConnectionConfig,
TAzureClientSecretsConnectionCredentials
} from "./azure-client-secrets-connection-types";
const generateClientAssertion = (
clientId: string,
tenantId: string,
privateKey: string,
certificate: string
): string => {
const tokenEndpoint = `https://login.microsoftonline.com/${tenantId}/oauth2/v2.0/token`;
const certBuffer = Buffer.from(
certificate
.replace(new RE2("-----BEGIN CERTIFICATE-----"), "")
.replace(new RE2("-----END CERTIFICATE-----"), "")
.replace(new RE2("\\s", "g"), ""),
"base64"
);
// thumbprint of the certificate is used for the jwt header
const thumbprint = crypto.nativeCrypto.createHash("sha1").update(certBuffer).digest("hex");
const x5t = Buffer.from(thumbprint, "hex").toString("base64url");
// JWT Header
const header = {
alg: "RS256",
typ: "JWT",
x5t
};
const now = Math.floor(Date.now() / 1000);
const payload = {
aud: tokenEndpoint,
exp: now + 600, // expire the assertion in 10 minutes (not the access access token TTL, but rather the assertion TTL itself)
iss: clientId,
jti: uuidv4(), // random ID for the JWT
nbf: now, // not before the jwt is valid
sub: clientId
};
// encode header and payload
const encodedHeader = Buffer.from(JSON.stringify(header)).toString("base64url");
const encodedPayload = Buffer.from(JSON.stringify(payload)).toString("base64url");
const signatureInput = `${encodedHeader}.${encodedPayload}`;
let keyObject: KeyObject;
try {
if (privateKey.includes("BEGIN PRIVATE KEY")) {
keyObject = crypto.nativeCrypto.createPrivateKey(privateKey);
} else {
// if user forgot to wrap in begin/end private key, decode and use as der format
keyObject = crypto.nativeCrypto.createPrivateKey({
key: Buffer.from(privateKey, "base64"),
format: "der",
type: "pkcs8"
});
}
} catch (error) {
throw new BadRequestError({
message: "Invalid private key format provided. Expected PEM format private key."
});
}
// sign with private key
const signer = crypto.nativeCrypto.createSign("RSA-SHA256");
signer.update(signatureInput);
signer.end();
const signature = signer.sign(keyObject, "base64url");
return `${signatureInput}.${signature}`;
};
export const getAzureClientSecretsConnectionListItem = () => {
const { INF_APP_CONNECTION_AZURE_CLIENT_SECRETS_CLIENT_ID } = getConfig();
@@ -30,7 +106,8 @@ export const getAzureClientSecretsConnectionListItem = () => {
app: AppConnection.AzureClientSecrets as const,
methods: Object.values(AzureClientSecretsConnectionMethod) as [
AzureClientSecretsConnectionMethod.OAuth,
AzureClientSecretsConnectionMethod.ClientSecret
AzureClientSecretsConnectionMethod.ClientSecret,
AzureClientSecretsConnectionMethod.Certificate
],
oauthClientId: INF_APP_CONNECTION_AZURE_CLIENT_SECRETS_CLIENT_ID
};
@@ -64,7 +141,7 @@ export const getAzureConnectionAccessToken = async (
const { refreshToken } = credentials;
const currentTime = Date.now();
switch (appConnection.method) {
case AzureClientSecretsConnectionMethod.OAuth:
case AzureClientSecretsConnectionMethod.OAuth: {
if (
!appCfg.INF_APP_CONNECTION_AZURE_CLIENT_SECRETS_CLIENT_ID ||
!appCfg.INF_APP_CONNECTION_AZURE_CLIENT_SECRETS_CLIENT_SECRET
@@ -101,7 +178,8 @@ export const getAzureConnectionAccessToken = async (
await appConnectionDAL.updateById(appConnection.id, { encryptedCredentials });
return data.access_token;
case AzureClientSecretsConnectionMethod.ClientSecret:
}
case AzureClientSecretsConnectionMethod.ClientSecret: {
const accessTokenCredentials = (await decryptAppConnectionCredentials({
orgId: appConnection.orgId,
projectId: appConnection.projectId,
@@ -139,6 +217,50 @@ export const getAzureConnectionAccessToken = async (
await appConnectionDAL.updateById(appConnection.id, { encryptedCredentials: encryptedClientCredentials });
return clientData.access_token;
}
case AzureClientSecretsConnectionMethod.Certificate: {
const accessTokenCredentials = (await decryptAppConnectionCredentials({
orgId: appConnection.orgId,
projectId: appConnection.projectId,
kmsService,
encryptedCredentials: appConnection.encryptedCredentials
})) as TAzureClientSecretsConnectionCertificateCredentials;
const { accessToken, expiresAt, clientId, tenantId, certificateBody, privateKey } = accessTokenCredentials;
if (accessToken && expiresAt && expiresAt > currentTime + 300000) {
return accessToken;
}
const clientAssertion = generateClientAssertion(clientId, tenantId, privateKey, certificateBody);
const { data: clientData } = await request.post<ExchangeCodeAzureResponse>(
IntegrationUrls.AZURE_TOKEN_URL.replace("common", tenantId || "common"),
new URLSearchParams({
grant_type: "client_credentials",
scope: `https://graph.microsoft.com/.default`,
client_id: clientId,
client_assertion_type: "urn:ietf:params:oauth:client-assertion-type:jwt-bearer",
client_assertion: clientAssertion
})
);
const updatedClientCredentials = {
...accessTokenCredentials,
accessToken: clientData.access_token,
expiresAt: currentTime + clientData.expires_in * 1000
};
const encryptedClientCredentials = await encryptAppConnectionCredentials({
credentials: updatedClientCredentials,
orgId: appConnection.orgId,
projectId: appConnection.projectId,
kmsService
});
await appConnectionDAL.updateById(appConnection.id, { encryptedCredentials: encryptedClientCredentials });
return clientData.access_token;
}
default:
throw new InternalServerError({
message: `Unhandled Azure connection method: ${appConnection.method as AzureClientSecretsConnectionMethod}`
@@ -156,7 +278,7 @@ export const validateAzureClientSecretsConnectionCredentials = async (config: TA
} = getConfig();
switch (method) {
case AzureClientSecretsConnectionMethod.OAuth:
case AzureClientSecretsConnectionMethod.OAuth: {
if (!SITE_URL) {
throw new InternalServerError({ message: "SITE_URL env var is required to complete Azure OAuth flow" });
}
@@ -221,8 +343,9 @@ export const validateAzureClientSecretsConnectionCredentials = async (config: TA
refreshToken: tokenResp.data.refresh_token,
expiresAt: Date.now() + tokenResp.data.expires_in * 1000
};
}
case AzureClientSecretsConnectionMethod.ClientSecret:
case AzureClientSecretsConnectionMethod.ClientSecret: {
const { tenantId, clientId, clientSecret } = inputCredentials;
try {
const { data: clientData } = await request.post<ExchangeCodeAzureResponse>(
@@ -255,6 +378,57 @@ export const validateAzureClientSecretsConnectionCredentials = async (config: TA
});
}
}
}
case AzureClientSecretsConnectionMethod.Certificate: {
const { tenantId, certificateBody, privateKey, clientId } = inputCredentials;
try {
const clientAssertion = generateClientAssertion(clientId, tenantId, privateKey, certificateBody);
const tokenEndpoint = `https://login.microsoftonline.com/${tenantId}/oauth2/v2.0/token`;
const params = new URLSearchParams({
client_id: clientId,
client_assertion_type: "urn:ietf:params:oauth:client-assertion-type:jwt-bearer",
client_assertion: clientAssertion,
scope: "https://graph.microsoft.com/.default",
grant_type: "client_credentials"
});
const response = await request.post<ExchangeCodeAzureResponse>(tokenEndpoint, params.toString(), {
headers: {
"Content-Type": "application/x-www-form-urlencoded"
}
});
return {
tenantId,
clientId,
certificateBody,
privateKey,
accessToken: response.data.access_token,
expiresAt: Date.now() + response.data.expires_in * 1000
};
} catch (e: unknown) {
if (e instanceof AxiosError) {
throw new BadRequestError({
message: `Failed to get access token: ${
(e?.response?.data as { error_description?: string })?.error_description || "Unknown error"
}`
});
} else if (e instanceof BadRequestError) {
throw e;
} else {
logger.error(
e,
"validateAzureClientSecretsConnectionCredentials: Failed to get access token using certificate authentication"
);
throw new InternalServerError({
message: "Failed to get access token"
});
}
}
}
default:
throw new InternalServerError({
message: `Unhandled Azure connection method: ${method as AzureClientSecretsConnectionMethod}`
@@ -48,6 +48,31 @@ export const AzureClientSecretsConnectionClientSecretInputCredentialsSchema = z.
.describe(AppConnections.CREDENTIALS.AZURE_CLIENT_SECRETS.tenantId)
});
export const AzureClientSecretsConnectionCertificateInputCredentialsSchema = z.object({
tenantId: z
.string()
.uuid()
.trim()
.min(1, "Tenant ID required")
.describe(AppConnections.CREDENTIALS.AZURE_CLIENT_SECRETS.tenantId),
clientId: z
.string()
.uuid()
.trim()
.min(1, "Client ID required")
.describe(AppConnections.CREDENTIALS.AZURE_CLIENT_SECRETS.clientId),
certificateBody: z
.string()
.trim()
.min(1, "Certificate body required")
.describe(AppConnections.CREDENTIALS.AZURE_CLIENT_SECRETS.certificateBody),
privateKey: z
.string()
.trim()
.min(1, "Private Key required")
.describe(AppConnections.CREDENTIALS.AZURE_CLIENT_SECRETS.privateKey)
});
export const AzureClientSecretsConnectionClientSecretOutputCredentialsSchema = z.object({
clientId: z.string(),
clientSecret: z.string(),
@@ -56,6 +81,15 @@ export const AzureClientSecretsConnectionClientSecretOutputCredentialsSchema = z
expiresAt: z.number()
});
export const AzureClientSecretsConnectionCertificateOutputCredentialsSchema = z.object({
clientId: z.string(),
tenantId: z.string(),
certificateBody: z.string(),
privateKey: z.string(),
accessToken: z.string(),
expiresAt: z.number()
});
export const ValidateAzureClientSecretsConnectionCredentialsSchema = z.discriminatedUnion("method", [
z.object({
method: z
@@ -72,6 +106,14 @@ export const ValidateAzureClientSecretsConnectionCredentialsSchema = z.discrimin
credentials: AzureClientSecretsConnectionClientSecretInputCredentialsSchema.describe(
AppConnections.CREATE(AppConnection.AzureClientSecrets).credentials
)
}),
z.object({
method: z
.literal(AzureClientSecretsConnectionMethod.Certificate)
.describe(AppConnections.CREATE(AppConnection.AzureClientSecrets).method),
credentials: AzureClientSecretsConnectionCertificateInputCredentialsSchema.describe(
AppConnections.CREATE(AppConnection.AzureClientSecrets).credentials
)
})
]);
@@ -84,7 +126,8 @@ export const UpdateAzureClientSecretsConnectionSchema = z
credentials: z
.union([
AzureClientSecretsConnectionOAuthInputCredentialsSchema,
AzureClientSecretsConnectionClientSecretInputCredentialsSchema
AzureClientSecretsConnectionClientSecretInputCredentialsSchema,
AzureClientSecretsConnectionCertificateInputCredentialsSchema
])
.optional()
.describe(AppConnections.UPDATE(AppConnection.AzureClientSecrets).credentials)
@@ -105,6 +148,10 @@ export const AzureClientSecretsConnectionSchema = z.intersection(
z.object({
method: z.literal(AzureClientSecretsConnectionMethod.ClientSecret),
credentials: AzureClientSecretsConnectionClientSecretOutputCredentialsSchema
}),
z.object({
method: z.literal(AzureClientSecretsConnectionMethod.Certificate),
credentials: AzureClientSecretsConnectionCertificateOutputCredentialsSchema
})
])
);
@@ -122,6 +169,13 @@ export const SanitizedAzureClientSecretsConnectionSchema = z.discriminatedUnion(
clientId: true,
tenantId: true
})
}),
BaseAzureClientSecretsConnectionSchema.extend({
method: z.literal(AzureClientSecretsConnectionMethod.Certificate),
credentials: AzureClientSecretsConnectionCertificateOutputCredentialsSchema.pick({
tenantId: true,
clientId: true
})
})
]);
@@ -4,6 +4,7 @@ import { DiscriminativePick } from "@app/lib/types";
import { AppConnection } from "../app-connection-enums";
import {
AzureClientSecretsConnectionCertificateOutputCredentialsSchema,
AzureClientSecretsConnectionClientSecretOutputCredentialsSchema,
AzureClientSecretsConnectionOAuthOutputCredentialsSchema,
AzureClientSecretsConnectionSchema,
@@ -35,6 +36,10 @@ export type TAzureClientSecretsConnectionClientSecretCredentials = z.infer<
typeof AzureClientSecretsConnectionClientSecretOutputCredentialsSchema
>;
export type TAzureClientSecretsConnectionCertificateCredentials = z.infer<
typeof AzureClientSecretsConnectionCertificateOutputCredentialsSchema
>;
export interface ExchangeCodeAzureResponse {
token_type: string;
scope: string;
@@ -0,0 +1,61 @@
import { Knex } from "knex";
import { TDbClient } from "@app/db";
import { TableName, TPkiAlertChannels, TPkiAlertChannelsInsert } from "@app/db/schemas";
import { DatabaseError } from "@app/lib/errors";
import { ormify, selectAllTableCols } from "@app/lib/knex";
export type TPkiAlertChannelDALFactory = ReturnType<typeof pkiAlertChannelDALFactory>;
export const pkiAlertChannelDALFactory = (db: TDbClient) => {
const pkiAlertChannelOrm = ormify(db, TableName.PkiAlertChannels);
const insertMany = async (data: TPkiAlertChannelsInsert[], tx?: Knex): Promise<TPkiAlertChannels[]> => {
try {
if (!data.length) return [];
const serializedData = data.map((item) => ({
...item,
config: item.config ? JSON.stringify(item.config) : null
}));
const res = await (tx || db)(TableName.PkiAlertChannels).insert(serializedData).returning("*");
return res as TPkiAlertChannels[];
} catch (error) {
throw new DatabaseError({ error, name: "InsertMany" });
}
};
const findByAlertId = async (alertId: string, tx?: Knex): Promise<TPkiAlertChannels[]> => {
try {
const channels = await (tx || db.replicaNode())(TableName.PkiAlertChannels)
.where(`${TableName.PkiAlertChannels}.alertId`, alertId)
.select(selectAllTableCols(TableName.PkiAlertChannels))
.orderBy(`${TableName.PkiAlertChannels}.createdAt`, "asc");
return channels as TPkiAlertChannels[];
} catch (error) {
throw new DatabaseError({ error, name: "FindByAlertId" });
}
};
const deleteByAlertId = async (alertId: string, tx?: Knex): Promise<number> => {
try {
const deletedCount = await (tx || db)(TableName.PkiAlertChannels)
.where(`${TableName.PkiAlertChannels}.alertId`, alertId)
.del();
return deletedCount;
} catch (error) {
throw new DatabaseError({ error, name: "DeleteByAlertId" });
}
};
return {
...pkiAlertChannelOrm,
insertMany,
findByAlertId,
deleteByAlertId
};
};
@@ -0,0 +1,110 @@
import { TDbClient } from "@app/db";
import { TableName, TPkiAlertHistory } from "@app/db/schemas";
import { DatabaseError } from "@app/lib/errors";
import { ormify, selectAllTableCols } from "@app/lib/knex";
export type TPkiAlertHistoryDALFactory = ReturnType<typeof pkiAlertHistoryDALFactory>;
export const pkiAlertHistoryDALFactory = (db: TDbClient) => {
const pkiAlertHistoryOrm = ormify(db, TableName.PkiAlertHistory);
const createWithCertificates = async (
alertId: string,
certificateIds: string[],
options?: {
hasNotificationSent?: boolean;
notificationError?: string;
}
): Promise<TPkiAlertHistory> => {
try {
return await db.transaction(async (tx) => {
const historyRecords = await tx(TableName.PkiAlertHistory)
.insert({
alertId,
hasNotificationSent: options?.hasNotificationSent || false,
notificationError: options?.notificationError
})
.returning("*");
const historyRecord = historyRecords[0];
if (certificateIds.length > 0) {
const certificateAssociations = certificateIds.map((certificateId) => ({
alertHistoryId: historyRecord.id,
certificateId
}));
await tx(TableName.PkiAlertHistoryCertificate).insert(certificateAssociations);
}
return historyRecord;
});
} catch (error) {
throw new DatabaseError({ error, name: "CreateWithCertificates" });
}
};
const findByAlertId = async (
alertId: string,
options?: {
limit?: number;
offset?: number;
}
): Promise<TPkiAlertHistory[]> => {
try {
let query = db
.replicaNode()
.select(selectAllTableCols(TableName.PkiAlertHistory))
.from(TableName.PkiAlertHistory)
.where(`${TableName.PkiAlertHistory}.alertId`, alertId)
.orderBy(`${TableName.PkiAlertHistory}.triggeredAt`, "desc");
if (options?.limit) {
query = query.limit(options.limit);
}
if (options?.offset) {
query = query.offset(options.offset);
}
const results = await query;
return results as TPkiAlertHistory[];
} catch (error) {
throw new DatabaseError({ error, name: "FindByAlertId" });
}
};
const findRecentlyAlertedCertificates = async (
alertId: string,
certificateIds: string[],
withinHours = 24
): Promise<string[]> => {
try {
if (certificateIds.length === 0) return [];
const cutoffDate = new Date();
cutoffDate.setHours(cutoffDate.getHours() - withinHours);
const results = (await db
.replicaNode()
.select("cert.certificateId")
.from(`${TableName.PkiAlertHistory} as hist`)
.join(`${TableName.PkiAlertHistoryCertificate} as cert`, "hist.id", "cert.alertHistoryId")
.where("hist.alertId", alertId)
.where("hist.hasNotificationSent", true)
.where("hist.triggeredAt", ">=", cutoffDate)
.whereIn("cert.certificateId", certificateIds)) as Array<{ certificateId: string }>;
return results.map((row) => row.certificateId);
} catch (error) {
throw new DatabaseError({ error, name: "FindRecentlyAlertedCertificates" });
}
};
return {
...pkiAlertHistoryOrm,
createWithCertificates,
findByAlertId,
findRecentlyAlertedCertificates
};
};
@@ -0,0 +1,556 @@
import { Knex } from "knex";
import { TDbClient } from "@app/db";
import { TableName, TPkiAlertsV2, TPkiAlertsV2Insert, TPkiAlertsV2Update } from "@app/db/schemas";
import { DatabaseError } from "@app/lib/errors";
import { ormify, selectAllTableCols } from "@app/lib/knex";
import {
applyCaFilters,
applyCertificateFilters,
requiresProfileJoin,
sanitizeLikeInput,
shouldIncludeCAs
} from "./pki-alert-v2-filter-utils";
import { CertificateOrigin, TCertificatePreview, TPkiFilterRule } from "./pki-alert-v2-types";
export type TPkiAlertV2DALFactory = ReturnType<typeof pkiAlertV2DALFactory>;
export const pkiAlertV2DALFactory = (db: TDbClient) => {
const pkiAlertV2Orm = ormify(db, TableName.PkiAlertsV2);
const create = async (data: TPkiAlertsV2Insert, tx?: Knex): Promise<TPkiAlertsV2> => {
try {
const serializedData = {
...data,
filters: data.filters ? JSON.stringify(data.filters) : null
};
const [res] = await (tx || db)(TableName.PkiAlertsV2).insert(serializedData).returning("*");
return res;
} catch (error) {
throw new DatabaseError({ error, name: "Create" });
}
};
const updateById = async (id: string, data: TPkiAlertsV2Update, tx?: Knex): Promise<TPkiAlertsV2> => {
try {
const serializedData: Record<string, unknown> = {
...data,
filters: data.filters !== undefined ? JSON.stringify(data.filters) : undefined
};
Object.keys(serializedData).forEach((key) => {
if (serializedData[key] === undefined) {
delete serializedData[key];
}
});
const [res] = await (tx || db)(TableName.PkiAlertsV2).where({ id }).update(serializedData).returning("*");
return res;
} catch (error) {
throw new DatabaseError({ error, name: "UpdateById" });
}
};
const findById = async (id: string, tx?: Knex): Promise<TPkiAlertsV2 | null> => {
try {
const [res] = await (tx || db.replicaNode())(TableName.PkiAlertsV2).where({ id }).select("*");
if (!res) return null;
return res;
} catch (error) {
throw new DatabaseError({ error, name: "FindById" });
}
};
type TChannelResult = {
id: string;
alertId: string;
channelType: string;
config: unknown;
enabled: boolean;
createdAt: Date;
updatedAt: Date;
};
type TAlertWithChannels = TPkiAlertsV2 & {
channels: TChannelResult[];
};
const findByIdWithChannels = async (alertId: string, tx?: Knex): Promise<TAlertWithChannels | null> => {
try {
const [alert] = (await (tx || db.replicaNode())
.select(selectAllTableCols(TableName.PkiAlertsV2))
.from(TableName.PkiAlertsV2)
.where(`${TableName.PkiAlertsV2}.id`, alertId)) as TPkiAlertsV2[];
if (!alert) return null;
const channels = (await (tx || db.replicaNode())
.select(selectAllTableCols(TableName.PkiAlertChannels))
.from(TableName.PkiAlertChannels)
.where(`${TableName.PkiAlertChannels}.alertId`, alertId)) as TChannelResult[];
return {
...alert,
channels: channels || []
} as TAlertWithChannels;
} catch (error) {
throw new DatabaseError({ error, name: "FindByIdWithChannels" });
}
};
const findByProjectIdWithCount = async (
projectId: string,
filters?: {
search?: string;
eventType?: string;
enabled?: boolean;
limit?: number;
offset?: number;
},
tx?: Knex
): Promise<{ alerts: TAlertWithChannels[]; total: number }> => {
try {
let countQuery = (tx || db.replicaNode())
.count("* as count")
.from(TableName.PkiAlertsV2)
.where(`${TableName.PkiAlertsV2}.projectId`, projectId);
if (filters?.search) {
countQuery = countQuery.whereILike(`${TableName.PkiAlertsV2}.name`, `%${sanitizeLikeInput(filters.search)}%`);
}
if (filters?.eventType) {
countQuery = countQuery.where(`${TableName.PkiAlertsV2}.eventType`, filters.eventType);
}
if (filters?.enabled !== undefined) {
countQuery = countQuery.where(`${TableName.PkiAlertsV2}.enabled`, filters.enabled);
}
let alertQuery = (tx || db.replicaNode())
.select(selectAllTableCols(TableName.PkiAlertsV2))
.from(TableName.PkiAlertsV2)
.where(`${TableName.PkiAlertsV2}.projectId`, projectId);
if (filters?.search) {
alertQuery = alertQuery.whereILike(`${TableName.PkiAlertsV2}.name`, `%${sanitizeLikeInput(filters.search)}%`);
}
if (filters?.eventType) {
alertQuery = alertQuery.where(`${TableName.PkiAlertsV2}.eventType`, filters.eventType);
}
if (filters?.enabled !== undefined) {
alertQuery = alertQuery.where(`${TableName.PkiAlertsV2}.enabled`, filters.enabled);
}
alertQuery = alertQuery.orderBy(`${TableName.PkiAlertsV2}.createdAt`, "desc");
if (filters?.limit) {
alertQuery = alertQuery.limit(filters.limit);
}
if (filters?.offset) {
alertQuery = alertQuery.offset(filters.offset);
}
const [countResult, alerts] = await Promise.all([countQuery, alertQuery]);
const total = parseInt((countResult[0] as { count: string }).count, 10);
const alertIds = (alerts as TPkiAlertsV2[]).map((alert) => alert.id);
const channels = (await (tx || db.replicaNode())
.select(selectAllTableCols(TableName.PkiAlertChannels))
.from(TableName.PkiAlertChannels)
.whereIn(`${TableName.PkiAlertChannels}.alertId`, alertIds)) as TChannelResult[];
const channelsByAlertId = channels.reduce(
(acc, channel) => {
if (!acc[channel.alertId]) {
acc[channel.alertId] = [];
}
acc[channel.alertId].push(channel);
return acc;
},
{} as Record<string, TChannelResult[]>
);
const alertsWithChannels: TAlertWithChannels[] = (alerts as TPkiAlertsV2[]).map((alert) => ({
...alert,
channels: channelsByAlertId[alert.id] || []
}));
return { alerts: alertsWithChannels, total };
} catch (error) {
throw new DatabaseError({ error, name: "FindByProjectIdWithCount" });
}
};
const findByProjectId = async (
projectId: string,
filters?: {
search?: string;
eventType?: string;
enabled?: boolean;
limit?: number;
offset?: number;
},
tx?: Knex
): Promise<TAlertWithChannels[]> => {
const result = await findByProjectIdWithCount(projectId, filters, tx);
return result.alerts;
};
const countByProjectId = async (
projectId: string,
filters?: {
search?: string;
eventType?: string;
enabled?: boolean;
},
tx?: Knex
): Promise<number> => {
try {
let query = (tx || db.replicaNode())
.count("* as count")
.from(TableName.PkiAlertsV2)
.where(`${TableName.PkiAlertsV2}.projectId`, projectId);
if (filters?.search) {
query = query.whereILike(`${TableName.PkiAlertsV2}.name`, `%${sanitizeLikeInput(filters.search)}%`);
}
if (filters?.eventType) {
query = query.where(`${TableName.PkiAlertsV2}.eventType`, filters.eventType);
}
if (filters?.enabled !== undefined) {
query = query.where(`${TableName.PkiAlertsV2}.enabled`, filters.enabled);
}
const result = await query;
return parseInt((result[0] as { count: string }).count, 10);
} catch (error) {
throw new DatabaseError({ error, name: "CountByProjectId" });
}
};
const getDistinctProjectIds = async (
filters?: {
enabled?: boolean;
},
tx?: Knex
): Promise<string[]> => {
try {
let query = (tx || db.replicaNode()).distinct(`${TableName.PkiAlertsV2}.projectId`).from(TableName.PkiAlertsV2);
if (filters?.enabled !== undefined) {
query = query.where(`${TableName.PkiAlertsV2}.enabled`, filters.enabled);
}
const result = await query;
return result.map((row: { projectId: string }) => row.projectId);
} catch (error) {
throw new DatabaseError({ error, name: "GetDistinctProjectIds" });
}
};
const findMatchingCertificates = async (
projectId: string,
filters: TPkiFilterRule[] = [],
options?: {
limit?: number;
offset?: number;
alertBefore?: string;
showFutureMatches?: boolean;
showCurrentMatches?: boolean;
showPreview?: boolean;
excludeAlerted?: boolean;
alertId?: string;
},
tx?: Knex
): Promise<{ certificates: TCertificatePreview[]; total: number }> => {
try {
const includeCAs = shouldIncludeCAs(filters);
const needsProfileJoin = requiresProfileJoin(filters);
const limit = options?.limit || 10;
const offset = options?.offset || 0;
let caTotalCount = 0;
let certTotalCount = 0;
if (includeCAs) {
let caCountQuery = (tx || db.replicaNode())
.count("* as count")
.from(TableName.CertificateAuthority)
.innerJoin(
`${TableName.InternalCertificateAuthority} as ica`,
`${TableName.CertificateAuthority}.id`,
`ica.caId`
);
caCountQuery = applyCaFilters(caCountQuery, filters, projectId) as typeof caCountQuery;
if (options?.alertBefore) {
if (options.showFutureMatches) {
caCountQuery = caCountQuery
.whereRaw(`ica."notAfter" > NOW() + ?::interval`, [options.alertBefore])
.whereRaw(`ica."notAfter" > NOW()`);
} else if (options.showCurrentMatches) {
caCountQuery = caCountQuery
.whereRaw(`ica."notAfter" > NOW()`)
.whereRaw(`ica."notAfter" <= NOW() + ?::interval`, [options.alertBefore]);
} else {
caCountQuery = caCountQuery
.whereRaw(`ica."notAfter" > NOW()`)
.whereRaw(`ica."notAfter" <= NOW() + ?::interval`, [options.alertBefore]);
}
}
const caCountResult = await caCountQuery;
caTotalCount = parseInt((caCountResult[0] as { count: string }).count, 10);
}
let certCountQuery = (tx || db.replicaNode()).count("* as count").from(TableName.Certificate);
certCountQuery = applyCertificateFilters(certCountQuery, filters, projectId) as typeof certCountQuery;
if (options?.showPreview) {
certCountQuery = certCountQuery
.whereRaw(`"${TableName.Certificate}"."notAfter" > NOW()`)
.whereNot(`${TableName.Certificate}.status`, "revoked");
} else if (options?.alertBefore) {
if (options.showFutureMatches) {
certCountQuery = certCountQuery
.whereRaw(`"${TableName.Certificate}"."notAfter" > NOW() + ?::interval`, [options.alertBefore])
.whereRaw(`"${TableName.Certificate}"."notAfter" > NOW()`)
.whereNot(`${TableName.Certificate}.status`, "revoked");
} else if (options.showCurrentMatches) {
certCountQuery = certCountQuery
.whereRaw(`"${TableName.Certificate}"."notAfter" > NOW()`)
.whereRaw(`"${TableName.Certificate}"."notAfter" <= NOW() + ?::interval`, [options.alertBefore])
.whereNot(`${TableName.Certificate}.status`, "revoked");
} else {
certCountQuery = certCountQuery
.whereRaw(`"${TableName.Certificate}"."notAfter" > NOW()`)
.whereRaw(`"${TableName.Certificate}"."notAfter" <= NOW() + ?::interval`, [options.alertBefore])
.whereNot(`${TableName.Certificate}.status`, "revoked");
}
}
if (options?.excludeAlerted && options?.alertId) {
certCountQuery = certCountQuery.whereNotExists(
(tx || db.replicaNode())
.select("*")
.from(TableName.PkiAlertHistory)
.join(
TableName.PkiAlertHistoryCertificate,
`${TableName.PkiAlertHistory}.id`,
`${TableName.PkiAlertHistoryCertificate}.alertHistoryId`
)
.where(`${TableName.PkiAlertHistory}.alertId`, options.alertId)
.whereRaw(`"${TableName.PkiAlertHistoryCertificate}"."certificateId" = "${TableName.Certificate}".id`)
);
}
const certCountResult = await certCountQuery;
certTotalCount = parseInt((certCountResult[0] as { count: string }).count, 10);
const totalCount = caTotalCount + certTotalCount;
let results: TCertificatePreview[] = [];
const fetchCertificates = async (certLimit: number, certOffset: number) => {
const selectColumns = [
`${TableName.Certificate}.id`,
`${TableName.Certificate}.serialNumber`,
`${TableName.Certificate}.commonName`,
`${TableName.Certificate}.altNames as san`,
`${TableName.Certificate}.notBefore`,
`${TableName.Certificate}.notAfter`,
`${TableName.Certificate}.status`,
`${TableName.Certificate}.profileId`,
`${TableName.Certificate}.pkiSubscriberId`
];
if (needsProfileJoin) {
selectColumns.push("profile.name as profileName");
}
let certificateQuery = (tx || db.replicaNode()).select(selectColumns).from(TableName.Certificate);
certificateQuery = applyCertificateFilters(certificateQuery, filters, projectId) as typeof certificateQuery;
if (options?.showPreview) {
certificateQuery = certificateQuery
.whereRaw(`"${TableName.Certificate}"."notAfter" > NOW()`)
.whereNot(`${TableName.Certificate}.status`, "revoked");
} else if (options?.alertBefore) {
if (options.showFutureMatches) {
certificateQuery = certificateQuery
.whereRaw(`"${TableName.Certificate}"."notAfter" > NOW() + ?::interval`, [options.alertBefore])
.whereRaw(`"${TableName.Certificate}"."notAfter" > NOW()`)
.whereNot(`${TableName.Certificate}.status`, "revoked");
} else if (options.showCurrentMatches) {
certificateQuery = certificateQuery
.whereRaw(`"${TableName.Certificate}"."notAfter" > NOW()`)
.whereRaw(`"${TableName.Certificate}"."notAfter" <= NOW() + ?::interval`, [options.alertBefore])
.whereNot(`${TableName.Certificate}.status`, "revoked");
} else {
certificateQuery = certificateQuery
.whereRaw(`"${TableName.Certificate}"."notAfter" > NOW()`)
.whereRaw(`"${TableName.Certificate}"."notAfter" <= NOW() + ?::interval`, [options.alertBefore])
.whereNot(`${TableName.Certificate}.status`, "revoked");
}
}
if (options?.excludeAlerted && options?.alertId) {
certificateQuery = certificateQuery.whereNotExists(
(tx || db.replicaNode())
.select("*")
.from(TableName.PkiAlertHistory)
.join(
TableName.PkiAlertHistoryCertificate,
`${TableName.PkiAlertHistory}.id`,
`${TableName.PkiAlertHistoryCertificate}.alertHistoryId`
)
.where(`${TableName.PkiAlertHistory}.alertId`, options.alertId)
.whereRaw(`"${TableName.PkiAlertHistoryCertificate}"."certificateId" = "${TableName.Certificate}".id`)
);
}
certificateQuery = certificateQuery
.orderBy(`${TableName.Certificate}.notAfter`, "asc")
.limit(certLimit)
.offset(certOffset);
const certificates = await certificateQuery;
const formattedCertificates: TCertificatePreview[] = (
certificates as Array<{
id: string;
serialNumber: string;
commonName: string;
san: string[] | null;
profileId: string | null;
pkiSubscriberId: string | null;
profileName?: string | null;
notBefore: Date;
notAfter: Date;
status: string;
}>
).map((cert) => {
let enrollmentType = CertificateOrigin.UNKNOWN;
if (cert.profileId) {
enrollmentType = CertificateOrigin.PROFILE;
} else if (cert.pkiSubscriberId) {
enrollmentType = CertificateOrigin.IMPORT;
}
return {
id: cert.id,
serialNumber: cert.serialNumber,
commonName: cert.commonName,
san: Array.isArray(cert.san) ? cert.san : [],
profileName: cert.profileName || null,
enrollmentType,
notBefore: cert.notBefore,
notAfter: cert.notAfter,
status: cert.status
};
});
results = [...results, ...formattedCertificates];
};
if (offset < caTotalCount) {
const caLimit = Math.min(limit, caTotalCount - offset);
const caOffset = offset;
let caQuery = (tx || db.replicaNode())
.select(
`${TableName.CertificateAuthority}.id`,
`ica.serialNumber`,
`ica.commonName`,
`ica.notBefore`,
`ica.notAfter`
)
.from(TableName.CertificateAuthority)
.innerJoin(
`${TableName.InternalCertificateAuthority} as ica`,
`${TableName.CertificateAuthority}.id`,
`ica.caId`
);
caQuery = applyCaFilters(caQuery, filters, projectId) as typeof caQuery;
if (options?.alertBefore) {
if (options.showFutureMatches) {
caQuery = caQuery
.whereRaw(`ica."notAfter" > NOW() + ?::interval`, [options.alertBefore])
.whereRaw(`ica."notAfter" > NOW()`);
} else {
caQuery = caQuery
.whereRaw(`ica."notAfter" > NOW()`)
.whereRaw(`ica."notAfter" <= NOW() + ?::interval`, [options.alertBefore]);
}
}
caQuery = caQuery.orderBy(`ica.notAfter`, "asc").limit(caLimit).offset(caOffset);
const cas = await caQuery;
const formattedCAs: TCertificatePreview[] = (
cas as Array<{
id: string;
serialNumber: string;
commonName: string;
notBefore: Date;
notAfter: Date;
}>
).map((ca) => ({
id: ca.id,
serialNumber: ca.serialNumber,
commonName: ca.commonName,
san: [],
profileName: null,
enrollmentType: CertificateOrigin.CA,
notBefore: ca.notBefore,
notAfter: ca.notAfter,
status: "active"
}));
results = [...results, ...formattedCAs];
const remainingLimit = limit - caLimit;
if (remainingLimit > 0 && certTotalCount > 0) {
const certOffset = 0;
await fetchCertificates(remainingLimit, certOffset);
}
} else {
const certOffset = offset - caTotalCount;
await fetchCertificates(limit, certOffset);
}
return {
certificates: results,
total: totalCount
};
} catch (error) {
throw new DatabaseError({ error, name: "FindMatchingCertificates" });
}
};
return {
...pkiAlertV2Orm,
create,
updateById,
findById,
findByIdWithChannels,
findByProjectId,
findByProjectIdWithCount,
countByProjectId,
getDistinctProjectIds,
findMatchingCertificates
};
};
@@ -0,0 +1,358 @@
import { Knex } from "knex";
import RE2 from "re2";
import { TableName } from "@app/db/schemas";
import { logger } from "@app/lib/logger";
import { PkiFilterField, PkiFilterOperator, TPkiFilterRule } from "./pki-alert-v2-types";
export const sanitizeLikeInput = (input: string): string => {
const allowedCharsRegex = new RE2("^[a-zA-Z0-9\\s\\-_\\.@\\*]+$");
if (!allowedCharsRegex.test(input)) {
throw new Error(
"Invalid characters in input. Only alphanumeric characters, spaces, hyphens, underscores, dots, @ and * are allowed."
);
}
const backslashRegex = new RE2("\\\\", "g");
const percentRegex = new RE2("%", "g");
const underscoreRegex = new RE2("_", "g");
const quoteRegex = new RE2("'", "g");
return input
.replace(backslashRegex, "\\\\\\\\")
.replace(percentRegex, "\\%")
.replace(underscoreRegex, "\\_")
.replace(quoteRegex, "''");
};
export const parseTimeToPostgresInterval = (duration: string): string => {
if (duration.length > 32) {
throw new Error(`Invalid duration format: ${duration}. Use format like '30d', '1w', '3m', '1y'`);
}
const durationRegex = new RE2("^(\\d+)([dwmy])$");
const match = durationRegex.exec(duration);
if (!match) {
throw new Error(`Invalid duration format: ${duration}. Use format like '30d', '1w', '3m', '1y'`);
}
const [, value, unit] = match;
const amount = parseInt(value, 10);
if (amount <= 0 || amount > 9999) {
throw new Error(`Duration value out of range: ${duration}. Must be between 1 and 9999.`);
}
const unitMap = {
d: "days",
w: "weeks",
m: "months",
y: "years"
};
return `${amount} ${unitMap[unit as keyof typeof unitMap]}`;
};
export const parseTimeToDays = (timeStr: string): number => {
const alertBeforeRegex = new RE2("^(\\d+)([dwmy])$");
const match = alertBeforeRegex.exec(timeStr);
if (!match) {
return 0;
}
const [, value, unit] = match;
const amount = parseInt(value, 10);
if (amount <= 0 || amount > 9999) {
return 0;
}
switch (unit) {
case "d":
return amount;
case "w":
return amount * 7;
case "m":
return amount * 30;
case "y":
return amount * 365;
default:
return 0;
}
};
const applyProfileNameFilter = (query: Knex.QueryBuilder, filter: TPkiFilterRule): Knex.QueryBuilder => {
const { value } = filter;
switch (filter.operator) {
case PkiFilterOperator.EQUALS:
return query.where("profile.slug", value as string);
case PkiFilterOperator.MATCHES:
if (Array.isArray(value)) {
return query.whereIn("profile.slug", value);
}
return query.whereILike("profile.slug", `%${sanitizeLikeInput(String(value))}%`);
case PkiFilterOperator.CONTAINS:
if (Array.isArray(value)) {
return query.where((builder) => {
value.forEach((v, index) => {
const sanitizedValue = sanitizeLikeInput(String(v));
if (index === 0) {
void builder.whereILike("profile.slug", `%${sanitizedValue}%`);
} else {
void builder.orWhereILike("profile.slug", `%${sanitizedValue}%`);
}
});
});
}
return query.whereILike("profile.slug", `%${sanitizeLikeInput(String(value))}%`);
case PkiFilterOperator.STARTS_WITH:
return query.whereILike("profile.slug", `${sanitizeLikeInput(String(value))}%`);
case PkiFilterOperator.ENDS_WITH:
return query.whereILike("profile.slug", `%${sanitizeLikeInput(String(value))}`);
default:
logger.warn(`Unsupported operator for profile_name: ${String(filter.operator)}`);
return query;
}
};
const applyCommonNameFilter = (query: Knex.QueryBuilder, filter: TPkiFilterRule): Knex.QueryBuilder => {
const { value } = filter;
const columnName = `${TableName.Certificate}.commonName`;
switch (filter.operator) {
case PkiFilterOperator.EQUALS:
return query.where(columnName, value as string);
case PkiFilterOperator.MATCHES:
if (Array.isArray(value)) {
return query.whereIn(columnName, value);
}
return query.whereILike(columnName, `%${sanitizeLikeInput(String(value))}%`);
case PkiFilterOperator.CONTAINS:
if (Array.isArray(value)) {
return query.where((builder) => {
value.forEach((v, index) => {
const sanitizedValue = sanitizeLikeInput(String(v));
if (index === 0) {
void builder.whereILike(columnName, `%${sanitizedValue}%`);
} else {
void builder.orWhereILike(columnName, `%${sanitizedValue}%`);
}
});
});
}
return query.whereILike(columnName, `%${sanitizeLikeInput(String(value))}%`);
case PkiFilterOperator.STARTS_WITH:
return query.whereILike(columnName, `${sanitizeLikeInput(String(value))}%`);
case PkiFilterOperator.ENDS_WITH:
return query.whereILike(columnName, `%${sanitizeLikeInput(String(value))}`);
default:
logger.warn(`Unsupported operator for common_name: ${String(filter.operator)}`);
return query;
}
};
const applySanFilter = (query: Knex.QueryBuilder, filter: TPkiFilterRule): Knex.QueryBuilder => {
const { value } = filter;
const columnName = `${TableName.Certificate}.altNames`;
switch (filter.operator) {
case PkiFilterOperator.EQUALS:
return query.whereJsonSupersetOf(columnName, [value as string]);
case PkiFilterOperator.MATCHES:
if (Array.isArray(value)) {
return query.where((builder) => {
value.forEach((v, index) => {
const sanitizedValue = `%"${String(v)}"%`;
if (index === 0) {
void builder.whereRaw(`??."altNames"::text ILIKE ?`, [TableName.Certificate, sanitizedValue]);
} else {
void builder.orWhereRaw(`??."altNames"::text ILIKE ?`, [TableName.Certificate, sanitizedValue]);
}
});
});
}
{
const sanitizedValue = `%"${String(value)}"%`;
return query.whereRaw(`??."altNames"::text ILIKE ?`, [TableName.Certificate, sanitizedValue]);
}
case PkiFilterOperator.CONTAINS:
return applySanFilter(query, { ...filter, operator: PkiFilterOperator.MATCHES });
case PkiFilterOperator.STARTS_WITH: {
const startsWithValue = `%"${String(value)}%`;
return query.whereRaw(`??."altNames"::text ILIKE ?`, [TableName.Certificate, startsWithValue]);
}
case PkiFilterOperator.ENDS_WITH: {
const endsWithValue = `%${String(value)}"%`;
return query.whereRaw(`??."altNames"::text ILIKE ?`, [TableName.Certificate, endsWithValue]);
}
default:
logger.warn(`Unsupported operator for SAN: ${String(filter.operator)}`);
return query;
}
};
export const shouldIncludeCAs = (filters: TPkiFilterRule[]): boolean => {
return filters.some((filter) => filter.field === PkiFilterField.INCLUDE_CAS && filter.value === true);
};
const applyCaCommonNameFilter = (query: Knex.QueryBuilder, filter: TPkiFilterRule): Knex.QueryBuilder => {
const { value } = filter;
const columnName = "ica.commonName";
switch (filter.operator) {
case PkiFilterOperator.EQUALS:
return query.where(columnName, value as string);
case PkiFilterOperator.MATCHES:
if (Array.isArray(value)) {
return query.whereIn(columnName, value);
}
return query.whereILike(columnName, `%${sanitizeLikeInput(String(value))}%`);
case PkiFilterOperator.CONTAINS:
if (Array.isArray(value)) {
return query.where((builder) => {
value.forEach((v, index) => {
if (index === 0) {
void builder.whereILike(columnName, `%${sanitizeLikeInput(String(v))}%`);
} else {
void builder.orWhereILike(columnName, `%${sanitizeLikeInput(String(v))}%`);
}
});
});
}
return query.whereILike(columnName, `%${sanitizeLikeInput(String(value))}%`);
case PkiFilterOperator.STARTS_WITH:
return query.whereILike(columnName, `${sanitizeLikeInput(String(value))}%`);
case PkiFilterOperator.ENDS_WITH:
return query.whereILike(columnName, `%${sanitizeLikeInput(String(value))}`);
default:
logger.warn(`Unsupported operator for CA common_name: ${String(filter.operator)}`);
return query;
}
};
export const applyCaFilters = (
query: Knex.QueryBuilder,
filters: TPkiFilterRule[],
projectId: string
): Knex.QueryBuilder => {
let filteredQuery = query.where(`${TableName.CertificateAuthority}.projectId`, projectId).whereNotNull("ica.caId"); // Only include CAs that have internal CA data
filters.forEach((filter) => {
switch (filter.field) {
case PkiFilterField.COMMON_NAME:
filteredQuery = applyCaCommonNameFilter(filteredQuery, filter);
break;
default:
break;
}
});
return filteredQuery;
};
export const validateFilterRules = (filters: TPkiFilterRule[]): void => {
for (const filter of filters) {
if (!Object.values(PkiFilterField).includes(filter.field)) {
throw new Error(`Invalid filter field: ${filter.field}`);
}
if (!Object.values(PkiFilterOperator).includes(filter.operator)) {
throw new Error(`Invalid filter operator: ${filter.operator}`);
}
switch (filter.field) {
case PkiFilterField.INCLUDE_CAS:
if (typeof filter.value !== "boolean") {
throw new Error("include_cas filter value must be boolean");
}
break;
case PkiFilterField.PROFILE_NAME:
case PkiFilterField.COMMON_NAME:
case PkiFilterField.SAN:
if (filter.operator === PkiFilterOperator.CONTAINS || filter.operator === PkiFilterOperator.MATCHES) {
if (!Array.isArray(filter.value) && typeof filter.value !== "string") {
throw new Error(
`${filter.field} filter value must be string or array of strings for ${filter.operator} operator`
);
}
} else if (typeof filter.value !== "string") {
throw new Error(`${filter.field} filter value must be string for ${filter.operator} operator`);
}
break;
default:
break;
}
}
};
export const requiresProfileJoin = (filters: TPkiFilterRule[]): boolean => {
return filters.some((filter) => filter.field === PkiFilterField.PROFILE_NAME);
};
export const applyCertificateFilters = (
query: Knex.QueryBuilder,
filters: TPkiFilterRule[],
projectId: string
): Knex.QueryBuilder => {
let filteredQuery = query.where(`${TableName.Certificate}.projectId`, projectId);
const needsProfileJoin = requiresProfileJoin(filters);
if (needsProfileJoin) {
filteredQuery = filteredQuery.leftJoin(
`${TableName.PkiCertificateProfile} as profile`,
`${TableName.Certificate}.profileId`,
"profile.id"
);
}
filters.forEach((filter) => {
switch (filter.field) {
case PkiFilterField.PROFILE_NAME:
filteredQuery = applyProfileNameFilter(filteredQuery, filter);
break;
case PkiFilterField.COMMON_NAME:
filteredQuery = applyCommonNameFilter(filteredQuery, filter);
break;
case PkiFilterField.SAN:
filteredQuery = applySanFilter(filteredQuery, filter);
break;
case PkiFilterField.INCLUDE_CAS:
break;
default:
logger.warn(`Unknown filter field: ${String(filter.field)}`);
break;
}
});
return filteredQuery;
};
@@ -0,0 +1,220 @@
/* eslint-disable no-await-in-loop */
import { getConfig } from "@app/lib/config/env";
import { logger } from "@app/lib/logger";
import { QueueJobs, QueueName, TQueueServiceFactory } from "@app/queue";
import { TPkiAlertHistoryDALFactory } from "./pki-alert-history-dal";
import { TPkiAlertV2DALFactory } from "./pki-alert-v2-dal";
import { parseTimeToDays, parseTimeToPostgresInterval } from "./pki-alert-v2-filter-utils";
import { TPkiAlertV2ServiceFactory } from "./pki-alert-v2-service";
import { CertificateOrigin, PkiAlertEventType, TPkiFilterRule } from "./pki-alert-v2-types";
type TPkiAlertV2QueueServiceFactoryDep = {
queueService: TQueueServiceFactory;
pkiAlertV2Service: Pick<TPkiAlertV2ServiceFactory, "sendAlertNotifications">;
pkiAlertV2DAL: Pick<TPkiAlertV2DALFactory, "findByProjectId" | "findMatchingCertificates" | "getDistinctProjectIds">;
pkiAlertHistoryDAL: Pick<TPkiAlertHistoryDALFactory, "findRecentlyAlertedCertificates">;
};
export type TPkiAlertV2QueueServiceFactory = ReturnType<typeof pkiAlertV2QueueServiceFactory>;
export const pkiAlertV2QueueServiceFactory = ({
queueService,
pkiAlertV2Service,
pkiAlertV2DAL,
pkiAlertHistoryDAL
}: TPkiAlertV2QueueServiceFactoryDep) => {
const appCfg = getConfig();
const calculateDeduplicationWindow = (alertBefore: string): number => {
const alertDays = parseTimeToDays(alertBefore);
if (alertDays === 0) {
return 24;
}
if (alertDays <= 1) {
return 8;
}
if (alertDays <= 7) {
return 24;
}
if (alertDays <= 30) {
return 48;
}
if (alertDays <= 90) {
return 168;
}
return 720;
};
const getAllProjectsWithAlerts = async (): Promise<string[]> => {
try {
const projectIds = await pkiAlertV2DAL.getDistinctProjectIds({ enabled: true });
logger.info(`Found ${projectIds.length} projects with PKI alerts`);
return projectIds;
} catch (error) {
logger.error(error, "Failed to get projects with alerts");
return [];
}
};
const evaluateAlert = async (
alert: {
id: string;
name: string;
eventType: string;
alertBefore: string;
filters: TPkiFilterRule[];
},
projectId: string
): Promise<{ shouldNotify: boolean; certificateIds: string[] }> => {
if (alert.eventType !== PkiAlertEventType.EXPIRATION) {
return { shouldNotify: false, certificateIds: [] };
}
try {
const result = await pkiAlertV2DAL.findMatchingCertificates(projectId, alert.filters, {
limit: 1000,
alertBefore: parseTimeToPostgresInterval(alert.alertBefore),
showCurrentMatches: true
});
if (result.certificates.length === 0) {
return { shouldNotify: false, certificateIds: [] };
}
const allCertificateIds = result.certificates
.filter((cert) => cert.enrollmentType !== CertificateOrigin.CA)
.map((cert) => cert.id);
const deduplicationHours = calculateDeduplicationWindow(alert.alertBefore);
const recentlyAlertedIds = await pkiAlertHistoryDAL.findRecentlyAlertedCertificates(
alert.id,
allCertificateIds,
deduplicationHours
);
const certificateIds = allCertificateIds.filter((certId) => !recentlyAlertedIds.includes(certId));
if (certificateIds.length === 0) {
logger.debug(
`All ${allCertificateIds.length} matching certificates for alert ${alert.id} were already alerted within the last ${deduplicationHours} hours`
);
return { shouldNotify: false, certificateIds: [] };
}
logger.debug(
`Alert ${alert.id}: Found ${allCertificateIds.length} expiring certificates, ${recentlyAlertedIds.length} already alerted recently, ${certificateIds.length} new to alert`
);
return {
shouldNotify: true,
certificateIds
};
} catch (error) {
logger.error(error, `Failed to evaluate alert ${alert.id}`);
return { shouldNotify: false, certificateIds: [] };
}
};
const processProjectAlerts = async (
projectId: string
): Promise<{ alertsProcessed: number; notificationsSent: number }> => {
logger.info(`Processing alerts for project: ${projectId}`);
const alerts = await pkiAlertV2DAL.findByProjectId(projectId, {
enabled: true,
limit: 1000
});
let alertsProcessed = 0;
let notificationsSent = 0;
for (const alert of alerts) {
const typedAlert = alert as {
id: string;
name: string;
eventType: string;
alertBefore: string;
filters: TPkiFilterRule[];
};
try {
const { shouldNotify, certificateIds } = await evaluateAlert(typedAlert, projectId);
if (shouldNotify && certificateIds.length > 0) {
await pkiAlertV2Service.sendAlertNotifications(typedAlert.id, certificateIds);
notificationsSent += 1;
logger.info(
`Sent notification for alert ${typedAlert.id} (${typedAlert.name}) with ${certificateIds.length} certificates`
);
}
alertsProcessed += 1;
} catch (error) {
logger.error(error, `Failed to process alert ${typedAlert.id} (${typedAlert.name})`);
}
}
logger.info(
`Completed processing ${alertsProcessed} alerts for project ${projectId}, sent ${notificationsSent} notifications`
);
return { alertsProcessed, notificationsSent };
};
const processDailyAlerts = async () => {
logger.info("Starting daily PKI alert processing...");
const allProjects = await getAllProjectsWithAlerts();
let totalAlertsProcessed = 0;
let totalNotificationsSent = 0;
for (const projectId of allProjects) {
try {
const { alertsProcessed, notificationsSent } = await processProjectAlerts(projectId);
totalAlertsProcessed += alertsProcessed;
totalNotificationsSent += notificationsSent;
} catch (error) {
logger.error(error, `Failed to process alerts for project ${projectId}`);
}
}
logger.info(
`Daily PKI alert processing completed. Processed ${totalAlertsProcessed} alerts, sent ${totalNotificationsSent} notifications.`
);
};
const init = async () => {
if (appCfg.isSecondaryInstance) {
return;
}
await queueService.startPg<QueueName.DailyPkiAlertV2Processing>(
QueueJobs.DailyPkiAlertV2Processing,
async () => {
try {
logger.info(`${QueueJobs.DailyPkiAlertV2Processing}: queue task started`);
await processDailyAlerts();
logger.info(`${QueueJobs.DailyPkiAlertV2Processing}: queue task completed successfully`);
} catch (error) {
logger.error(error, `${QueueJobs.DailyPkiAlertV2Processing}: queue task failed`);
throw error;
}
},
{
batchSize: 1,
workerCount: 1,
pollingIntervalSeconds: 60
}
);
await queueService.schedulePg(QueueJobs.DailyPkiAlertV2Processing, "0 0 * * *", undefined, { tz: "UTC" });
};
return {
init
};
};
@@ -0,0 +1,507 @@
import { ForbiddenError } from "@casl/ability";
import { ActionProjectType } from "@app/db/schemas";
import { TPermissionServiceFactory } from "@app/ee/services/permission/permission-service-types";
import { ProjectPermissionActions, ProjectPermissionSub } from "@app/ee/services/permission/project-permission";
import { BadRequestError, NotFoundError } from "@app/lib/errors";
import { logger } from "@app/lib/logger";
import { SmtpTemplates, TSmtpService } from "@app/services/smtp/smtp-service";
import { TPkiAlertChannelDALFactory } from "./pki-alert-channel-dal";
import { TPkiAlertHistoryDALFactory } from "./pki-alert-history-dal";
import { TPkiAlertV2DALFactory } from "./pki-alert-v2-dal";
import { parseTimeToDays, parseTimeToPostgresInterval } from "./pki-alert-v2-filter-utils";
import {
CertificateOrigin,
PkiAlertChannelType,
PkiAlertEventType,
TAlertV2Response,
TChannelConfig,
TCreateAlertV2DTO,
TDeleteAlertV2DTO,
TEmailChannelConfig,
TGetAlertV2DTO,
TListAlertsV2DTO,
TListAlertsV2Response,
TListCurrentMatchingCertificatesDTO,
TListMatchingCertificatesDTO,
TListMatchingCertificatesResponse,
TPkiFilterRule,
TUpdateAlertV2DTO
} from "./pki-alert-v2-types";
type TPkiAlertV2ServiceFactoryDep = {
pkiAlertV2DAL: Pick<
TPkiAlertV2DALFactory,
| "create"
| "findById"
| "findByIdWithChannels"
| "updateById"
| "deleteById"
| "findByProjectId"
| "findByProjectIdWithCount"
| "countByProjectId"
| "findMatchingCertificates"
| "transaction"
>;
pkiAlertChannelDAL: Pick<TPkiAlertChannelDALFactory, "create" | "findByAlertId" | "deleteByAlertId" | "insertMany">;
pkiAlertHistoryDAL: Pick<TPkiAlertHistoryDALFactory, "createWithCertificates" | "findByAlertId">;
permissionService: Pick<TPermissionServiceFactory, "getProjectPermission">;
smtpService: Pick<TSmtpService, "sendMail">;
};
export type TPkiAlertV2ServiceFactory = ReturnType<typeof pkiAlertV2ServiceFactory>;
export const pkiAlertV2ServiceFactory = ({
pkiAlertV2DAL,
pkiAlertChannelDAL,
pkiAlertHistoryDAL,
permissionService,
smtpService
}: TPkiAlertV2ServiceFactoryDep) => {
type TAlertWithChannels = {
id: string;
name: string;
description: string;
eventType: string;
alertBefore: string;
filters: TPkiFilterRule[];
enabled: boolean;
projectId: string;
createdAt: Date;
updatedAt: Date;
channels?: Array<{
id: string;
channelType: string;
config: unknown;
enabled: boolean;
createdAt: Date;
updatedAt: Date;
}>;
};
const formatAlertResponse = (alert: TAlertWithChannels): TAlertV2Response => {
return {
id: alert.id,
name: alert.name,
description: alert.description,
eventType: alert.eventType as PkiAlertEventType,
alertBefore: alert.alertBefore,
filters: alert.filters,
enabled: alert.enabled,
projectId: alert.projectId,
channels: (alert.channels || []).map((channel) => ({
id: channel.id,
channelType: channel.channelType as PkiAlertChannelType,
config: channel.config as TChannelConfig,
enabled: channel.enabled,
createdAt: channel.createdAt,
updatedAt: channel.updatedAt
})),
createdAt: alert.createdAt,
updatedAt: alert.updatedAt
};
};
const createAlert = async ({
projectId,
name,
description,
eventType,
alertBefore,
filters,
enabled = true,
channels,
actorId,
actorAuthMethod,
actor,
actorOrgId
}: TCreateAlertV2DTO): Promise<TAlertV2Response> => {
const { permission } = await permissionService.getProjectPermission({
actor,
actorId,
projectId,
actorAuthMethod,
actorOrgId,
actionProjectType: ActionProjectType.CertificateManager
});
ForbiddenError.from(permission).throwUnlessCan(ProjectPermissionActions.Create, ProjectPermissionSub.PkiAlerts);
try {
parseTimeToPostgresInterval(alertBefore);
} catch (error) {
throw new BadRequestError({ message: "Invalid alertBefore format. Use format like '30d', '1w', '3m', '1y'" });
}
return pkiAlertV2DAL.transaction(async (tx) => {
const alert = await pkiAlertV2DAL.create(
{
projectId,
name,
description,
eventType,
alertBefore,
filters,
enabled
},
tx
);
const channelInserts = channels.map((channel) => ({
alertId: alert.id,
channelType: channel.channelType,
config: channel.config,
enabled: channel.enabled
}));
await pkiAlertChannelDAL.insertMany(channelInserts, tx);
const completeAlert = await pkiAlertV2DAL.findByIdWithChannels(alert.id, tx);
if (!completeAlert) {
throw new NotFoundError({ message: "Failed to retrieve created alert" });
}
return formatAlertResponse(completeAlert as TAlertWithChannels);
});
};
const getAlertById = async ({
alertId,
actorId,
actorAuthMethod,
actor,
actorOrgId
}: TGetAlertV2DTO): Promise<TAlertV2Response> => {
const alert = await pkiAlertV2DAL.findByIdWithChannels(alertId);
if (!alert) throw new NotFoundError({ message: `Alert with ID '${alertId}' not found` });
const { permission } = await permissionService.getProjectPermission({
actor,
actorId,
projectId: (alert as { projectId: string }).projectId,
actorAuthMethod,
actorOrgId,
actionProjectType: ActionProjectType.CertificateManager
});
ForbiddenError.from(permission).throwUnlessCan(ProjectPermissionActions.Read, ProjectPermissionSub.PkiAlerts);
return formatAlertResponse(alert as TAlertWithChannels);
};
const listAlerts = async ({
projectId,
search,
eventType,
enabled,
limit = 20,
offset = 0,
actorId,
actorAuthMethod,
actor,
actorOrgId
}: TListAlertsV2DTO): Promise<TListAlertsV2Response> => {
const { permission } = await permissionService.getProjectPermission({
actor,
actorId,
projectId,
actorAuthMethod,
actorOrgId,
actionProjectType: ActionProjectType.CertificateManager
});
ForbiddenError.from(permission).throwUnlessCan(ProjectPermissionActions.Read, ProjectPermissionSub.PkiAlerts);
const filters = { search, eventType, enabled, limit, offset };
const { alerts, total } = await pkiAlertV2DAL.findByProjectIdWithCount(projectId, filters);
return {
alerts: alerts.map((alert) => formatAlertResponse(alert as TAlertWithChannels)),
total
};
};
const updateAlert = async ({
alertId,
name,
description,
eventType,
alertBefore,
filters,
enabled,
channels,
actorId,
actorAuthMethod,
actor,
actorOrgId
}: TUpdateAlertV2DTO): Promise<TAlertV2Response> => {
let alert = await pkiAlertV2DAL.findById(alertId);
if (!alert) throw new NotFoundError({ message: `Alert with ID '${alertId}' not found` });
const { permission } = await permissionService.getProjectPermission({
actor,
actorId,
projectId: (alert as { projectId: string }).projectId,
actorAuthMethod,
actorOrgId,
actionProjectType: ActionProjectType.CertificateManager
});
ForbiddenError.from(permission).throwUnlessCan(ProjectPermissionActions.Edit, ProjectPermissionSub.PkiAlerts);
if (alertBefore) {
try {
parseTimeToPostgresInterval(alertBefore);
} catch (error) {
throw new BadRequestError({ message: "Invalid alertBefore format. Use format like '30d', '1w', '3m', '1y'" });
}
}
const updateData: {
name?: string;
description?: string;
eventType?: PkiAlertEventType;
alertBefore?: string;
filters?: TPkiFilterRule[];
enabled?: boolean;
} = {};
if (name !== undefined) updateData.name = name;
if (description !== undefined) updateData.description = description;
if (eventType !== undefined) updateData.eventType = eventType;
if (alertBefore !== undefined) updateData.alertBefore = alertBefore;
if (filters !== undefined) updateData.filters = filters;
if (enabled !== undefined) updateData.enabled = enabled;
return pkiAlertV2DAL.transaction(async (tx) => {
alert = await pkiAlertV2DAL.updateById(alertId, updateData, tx);
if (channels) {
await pkiAlertChannelDAL.deleteByAlertId(alertId, tx);
const channelInserts = channels.map((channel) => ({
alertId,
channelType: channel.channelType,
config: channel.config,
enabled: channel.enabled
}));
await pkiAlertChannelDAL.insertMany(channelInserts, tx);
}
const completeAlert = await pkiAlertV2DAL.findByIdWithChannels(alertId, tx);
if (!completeAlert) {
throw new NotFoundError({ message: "Failed to retrieve updated alert" });
}
return formatAlertResponse(completeAlert as TAlertWithChannels);
});
};
const deleteAlert = async ({
alertId,
actorId,
actorAuthMethod,
actor,
actorOrgId
}: TDeleteAlertV2DTO): Promise<TAlertV2Response> => {
const alert = await pkiAlertV2DAL.findByIdWithChannels(alertId);
if (!alert) throw new NotFoundError({ message: `Alert with ID '${alertId}' not found` });
const { permission } = await permissionService.getProjectPermission({
actor,
actorId,
projectId: (alert as { projectId: string }).projectId,
actorAuthMethod,
actorOrgId,
actionProjectType: ActionProjectType.CertificateManager
});
ForbiddenError.from(permission).throwUnlessCan(ProjectPermissionActions.Delete, ProjectPermissionSub.PkiAlerts);
const formattedAlert = formatAlertResponse(alert as TAlertWithChannels);
await pkiAlertV2DAL.deleteById(alertId);
return formattedAlert;
};
const listMatchingCertificates = async ({
alertId,
limit = 20,
offset = 0,
actorId,
actorAuthMethod,
actor,
actorOrgId
}: TListMatchingCertificatesDTO): Promise<TListMatchingCertificatesResponse> => {
const alert = await pkiAlertV2DAL.findById(alertId);
if (!alert) throw new NotFoundError({ message: `Alert with ID '${alertId}' not found` });
const { permission } = await permissionService.getProjectPermission({
actor,
actorId,
projectId: (alert as { projectId: string }).projectId,
actorAuthMethod,
actorOrgId,
actionProjectType: ActionProjectType.CertificateManager
});
ForbiddenError.from(permission).throwUnlessCan(ProjectPermissionActions.Read, ProjectPermissionSub.PkiAlerts);
const options: {
limit: number;
offset: number;
showPreview?: boolean;
excludeAlerted?: boolean;
alertId?: string;
} = {
limit,
offset,
showPreview: true,
excludeAlerted: (alert as { eventType: string }).eventType === PkiAlertEventType.EXPIRATION,
alertId
};
const result = await pkiAlertV2DAL.findMatchingCertificates(
(alert as { projectId: string }).projectId,
(alert as { filters: TPkiFilterRule[] }).filters,
options
);
return {
certificates: result.certificates,
total: result.total
};
};
const listCurrentMatchingCertificates = async ({
projectId,
filters,
alertBefore,
limit = 20,
offset = 0,
actorId,
actorAuthMethod,
actor,
actorOrgId
}: TListCurrentMatchingCertificatesDTO): Promise<TListMatchingCertificatesResponse> => {
const { permission } = await permissionService.getProjectPermission({
actor,
actorId,
projectId,
actorAuthMethod,
actorOrgId,
actionProjectType: ActionProjectType.CertificateManager
});
ForbiddenError.from(permission).throwUnlessCan(ProjectPermissionActions.Read, ProjectPermissionSub.PkiAlerts);
try {
parseTimeToPostgresInterval(alertBefore);
} catch (error) {
throw new BadRequestError({ message: "Invalid alertBefore format. Use format like '30d', '1w', '3m', '1y'" });
}
const options: {
limit: number;
offset: number;
showPreview?: boolean;
alertBefore?: string;
} = {
limit,
offset,
showPreview: true,
alertBefore: parseTimeToPostgresInterval(alertBefore)
};
const result = await pkiAlertV2DAL.findMatchingCertificates(projectId, filters, options);
return {
certificates: result.certificates,
total: result.total
};
};
const sendAlertNotifications = async (alertId: string, certificateIds: string[]) => {
const alert = await pkiAlertV2DAL.findByIdWithChannels(alertId);
if (!alert || !(alert as { enabled: boolean }).enabled) return;
const channels =
(alert as { channels?: Array<{ enabled: boolean; channelType: string; config: unknown }> }).channels?.filter(
(channel: { enabled: boolean; channelType: string; config: unknown }) => channel.enabled
) || [];
if (channels.length === 0) return;
const { certificates } = await pkiAlertV2DAL.findMatchingCertificates(
(alert as { projectId: string }).projectId,
(alert as { filters: TPkiFilterRule[] }).filters,
{
alertBefore: parseTimeToPostgresInterval((alert as { alertBefore: string }).alertBefore)
}
);
const matchingCertificates = certificates.filter(
(cert) => certificateIds.includes(cert.id) && cert.enrollmentType !== CertificateOrigin.CA
);
if (matchingCertificates.length === 0) return;
let hasNotificationSent = false;
let notificationError: string | undefined;
try {
const emailChannels = channels.filter(
(channel: { enabled: boolean; channelType: string; config: unknown }) =>
channel.channelType === PkiAlertChannelType.EMAIL
);
const alertBeforeDays = parseTimeToDays((alert as { alertBefore: string }).alertBefore);
const alertName = (alert as { name: string }).name;
const emailPromises = emailChannels.map((channel) => {
const config = channel.config as TEmailChannelConfig;
return smtpService.sendMail({
recipients: config.recipients,
subjectLine: `Infisical Certificate Alert - ${alertName}`,
substitutions: {
alertName,
alertBeforeDays,
projectId: (alert as { projectId: string }).projectId,
items: matchingCertificates.map((cert) => ({
type: "Certificate",
friendlyName: cert.commonName,
serialNumber: cert.serialNumber,
expiryDate: cert.notAfter.toLocaleDateString()
}))
},
template: SmtpTemplates.PkiExpirationAlert
});
});
await Promise.all(emailPromises);
hasNotificationSent = true;
} catch (error) {
notificationError = error instanceof Error ? error.message : "Unknown error occurred";
logger.error(error, `Failed to send notifications for alert ${alertId}`);
}
await pkiAlertHistoryDAL.createWithCertificates(alertId, certificateIds, {
hasNotificationSent,
notificationError
});
};
return {
createAlert,
getAlertById,
listAlerts,
updateAlert,
deleteAlert,
listMatchingCertificates,
listCurrentMatchingCertificates,
sendAlertNotifications
};
};
@@ -0,0 +1,205 @@
import RE2 from "re2";
import { z } from "zod";
import { TGenericPermission } from "@app/lib/types";
const createSecureNameValidator = () => {
// Validates name format: lowercase alphanumeric characters with optional hyphens
// Pattern: starts and ends with alphanumeric, allows hyphens between segments
// Examples: "my-alert", "alert1", "test-alert-2"
const nameRegex = new RE2("^[a-z0-9]+(?:-[a-z0-9]+)*$");
return (value: string) => nameRegex.test(value);
};
export const createSecureAlertBeforeValidator = () => {
// Validates alertBefore duration format: number followed by time unit
// Pattern: one or more digits followed by d(days), w(weeks), m(months), or y(years)
// Examples: "30d", "2w", "6m", "1y"
const alertBeforeRegex = new RE2("^\\d+[dwmy]$");
return (value: string) => {
if (value.length > 32) return false;
return alertBeforeRegex.test(value);
};
};
export enum PkiAlertEventType {
EXPIRATION = "expiration",
RENEWAL = "renewal",
ISSUANCE = "issuance",
REVOCATION = "revocation"
}
export enum PkiAlertChannelType {
EMAIL = "email",
WEBHOOK = "webhook",
SLACK = "slack"
}
export enum PkiFilterOperator {
EQUALS = "equals",
MATCHES = "matches",
CONTAINS = "contains",
STARTS_WITH = "starts_with",
ENDS_WITH = "ends_with"
}
export enum PkiFilterField {
PROFILE_NAME = "profile_name",
COMMON_NAME = "common_name",
SAN = "san",
INCLUDE_CAS = "include_cas"
}
export enum CertificateOrigin {
UNKNOWN = "unknown",
PROFILE = "profile",
IMPORT = "import",
CA = "ca"
}
export const PkiFilterRuleSchema = z.object({
field: z.nativeEnum(PkiFilterField),
operator: z.nativeEnum(PkiFilterOperator),
value: z.union([z.string(), z.array(z.string()), z.boolean()])
});
export type TPkiFilterRule = z.infer<typeof PkiFilterRuleSchema>;
export const PkiFiltersSchema = z.array(PkiFilterRuleSchema);
export type TPkiFilters = z.infer<typeof PkiFiltersSchema>;
export const EmailChannelConfigSchema = z.object({
recipients: z.array(z.string().email()).min(1).max(10)
});
export const WebhookChannelConfigSchema = z.object({
url: z.string().url(),
method: z.enum(["POST", "PUT"]).default("POST"),
headers: z.record(z.string()).optional()
});
export const SlackChannelConfigSchema = z.object({
webhookUrl: z.string().url(),
channel: z.string().optional(),
mentionUsers: z.array(z.string()).optional()
});
export const ChannelConfigSchema = z.union([
EmailChannelConfigSchema,
WebhookChannelConfigSchema,
SlackChannelConfigSchema
]);
export type TEmailChannelConfig = z.infer<typeof EmailChannelConfigSchema>;
export type TWebhookChannelConfig = z.infer<typeof WebhookChannelConfigSchema>;
export type TSlackChannelConfig = z.infer<typeof SlackChannelConfigSchema>;
export type TChannelConfig = z.infer<typeof ChannelConfigSchema>;
export const CreateChannelSchema = z.object({
channelType: z.nativeEnum(PkiAlertChannelType),
config: ChannelConfigSchema,
enabled: z.boolean().default(true)
});
export type TCreateChannel = z.infer<typeof CreateChannelSchema>;
export const CreatePkiAlertV2Schema = z.object({
name: z
.string()
.min(1)
.max(255)
.refine(createSecureNameValidator(), "Must be a valid name (lowercase, numbers, hyphens only)"),
description: z.string().max(1000).optional(),
eventType: z.nativeEnum(PkiAlertEventType),
alertBefore: z.string().refine(createSecureAlertBeforeValidator(), "Must be in format like '30d', '1w', '3m', '1y'"),
filters: PkiFiltersSchema,
enabled: z.boolean().default(true),
channels: z.array(CreateChannelSchema).min(1, "At least one channel is required")
});
export type TCreatePkiAlertV2 = z.infer<typeof CreatePkiAlertV2Schema>;
export const UpdatePkiAlertV2Schema = CreatePkiAlertV2Schema.partial();
export type TUpdatePkiAlertV2 = z.infer<typeof UpdatePkiAlertV2Schema>;
export type TCreateAlertV2DTO = TGenericPermission & {
projectId: string;
} & TCreatePkiAlertV2;
export type TUpdateAlertV2DTO = TGenericPermission & {
alertId: string;
} & TUpdatePkiAlertV2;
export type TGetAlertV2DTO = TGenericPermission & {
alertId: string;
};
export type TDeleteAlertV2DTO = TGenericPermission & {
alertId: string;
};
export type TListAlertsV2DTO = TGenericPermission & {
projectId: string;
search?: string;
eventType?: PkiAlertEventType;
enabled?: boolean;
limit?: number;
offset?: number;
};
export type TListMatchingCertificatesDTO = TGenericPermission & {
alertId: string;
limit?: number;
offset?: number;
};
export type TListCurrentMatchingCertificatesDTO = TGenericPermission & {
projectId: string;
filters: TPkiFilters;
alertBefore: string;
limit?: number;
offset?: number;
};
export type TCertificatePreview = {
id: string;
serialNumber: string;
commonName: string;
san: string[];
profileName: string | null;
enrollmentType: CertificateOrigin | null;
notBefore: Date;
notAfter: Date;
status: string;
};
export type TAlertV2Response = {
id: string;
name: string;
description: string | null;
eventType: PkiAlertEventType;
alertBefore: string;
filters: TPkiFilters;
enabled: boolean;
projectId: string;
channels: Array<{
id: string;
channelType: PkiAlertChannelType;
config: TChannelConfig;
enabled: boolean;
createdAt: Date;
updatedAt: Date;
}>;
createdAt: Date;
updatedAt: Date;
};
export type TListAlertsV2Response = {
alerts: TAlertV2Response[];
total: number;
};
export type TListMatchingCertificatesResponse = {
certificates: TCertificatePreview[];
total: number;
};
@@ -1,11 +1,12 @@
import { Heading, Hr, Section, Text } from "@react-email/components";
import React, { Fragment } from "react";
import { Heading, Section, Text } from "@react-email/components";
import { BaseButton } from "./BaseButton";
import { BaseEmailWrapper, BaseEmailWrapperProps } from "./BaseEmailWrapper";
interface PkiExpirationAlertTemplateProps extends Omit<BaseEmailWrapperProps, "title" | "preview" | "children"> {
alertName: string;
alertBeforeDays: number;
projectId: string;
items: { type: string; friendlyName: string; serialNumber: string; expiryDate: string }[];
}
@@ -13,44 +14,56 @@ export const PkiExpirationAlertTemplate = ({
alertName,
siteUrl,
alertBeforeDays,
projectId,
items
}: PkiExpirationAlertTemplateProps) => {
const formatDate = (dateStr: string) => {
try {
return new Date(dateStr).toLocaleDateString("en-US", {
year: "numeric",
month: "long",
day: "numeric"
});
} catch {
return dateStr;
}
};
const certificateText = items.length === 1 ? "certificate" : "certificates";
const daysText = alertBeforeDays === 1 ? "1 day" : `${alertBeforeDays} days`;
const message = `Alert ${alertName}: You have ${items.length === 1 ? "one" : items.length} ${certificateText} that will expire in ${daysText}.`;
return (
<BaseEmailWrapper
title="Infisical CA/Certificate Expiration Notice"
preview="One or more of your Infisical certificates is about to expire."
siteUrl={siteUrl}
>
<BaseEmailWrapper title="Certificate Expiration Notice" preview={message} siteUrl={siteUrl}>
<Heading className="text-black text-[18px] leading-[28px] text-center font-normal p-0 mx-0">
<strong>CA/Certificate Expiration Notice</strong>
<strong>Certificate Expiration Notice</strong>
</Heading>
<Section className="px-[24px] mt-[36px] pt-[12px] pb-[8px] border border-solid border-gray-200 rounded-md bg-gray-50">
<Text>Hello,</Text>
<Text className="text-black text-[14px] leading-[24px]">
This is an automated alert for <strong>{alertName}</strong> triggered for CAs/Certificates expiring in{" "}
<strong>{alertBeforeDays}</strong> days.
</Text>
<Text className="text-[14px] leading-[24px] mb-[4px]">
<strong>Expiring Items:</strong>
<Section className="px-[24px] mb-[28px] mt-[36px] pt-[12px] pb-[8px] border border-solid border-gray-200 rounded-md bg-gray-50">
<Text className="text-[14px]">
Alert <strong className="font-semibold">{alertName}</strong>: You have{" "}
{items.length === 1 ? "one" : items.length} {certificateText} that will expire in {daysText}.
</Text>
</Section>
<Section className="mb-[28px]">
<Text className="text-[14px] font-semibold mb-[12px]">Expiring certificates:</Text>
{items.map((item) => (
<Fragment key={item.serialNumber}>
<Hr className="mb-[16px]" />
<strong className="text-[14px]">{item.type}:</strong>
<Text className="text-[14px] my-[2px] leading-[24px]">{item.friendlyName}</Text>
<strong className="text-[14px]">Serial Number:</strong>
<Text className="text-[14px] my-[2px] leading-[24px]">{item.serialNumber}</Text>
<strong className="text-[14px]">Expires On:</strong>
<Text className="text-[14px] mt-[2px] mb-[16px] leading-[24px]">{item.expiryDate}</Text>
</Fragment>
<Section
key={item.serialNumber}
className="mb-[16px] p-[16px] border border-solid border-gray-200 rounded-md bg-gray-50"
>
<Text className="text-[14px] font-semibold m-0 mb-[4px]">{item.friendlyName}</Text>
<Text className="text-[12px] text-gray-600 m-0 mb-[4px]">Serial: {item.serialNumber}</Text>
<Text className="text-[12px] text-gray-600 m-0">Expires: {formatDate(item.expiryDate)}</Text>
</Section>
))}
<Hr />
<Text className="text-[14px] leading-[24px]">
Please take the necessary actions to renew these items before they expire.
</Text>
<Text className="text-[14px] leading-[24px]">
For more details, please log in to your Infisical account and check your PKI management section.
</Text>
</Section>
<Section className="text-center mt-[32px] mb-[16px]">
<BaseButton href={`${siteUrl}/projects/cert-management/${projectId}/policies`}>
View Certificate Alerts
</BaseButton>
</Section>
</BaseEmailWrapper>
);
@@ -59,11 +72,22 @@ export const PkiExpirationAlertTemplate = ({
export default PkiExpirationAlertTemplate;
PkiExpirationAlertTemplate.PreviewProps = {
alertBeforeDays: 5,
alertBeforeDays: 7,
items: [
{ type: "CA", friendlyName: "Example CA", serialNumber: "1234567890", expiryDate: "2032-01-01" },
{ type: "Certificate", friendlyName: "Example Certificate", serialNumber: "2345678901", expiryDate: "2032-01-01" }
{
type: "Certificate",
friendlyName: "api.production.company.com",
serialNumber: "4B:3E:2F:A1:D6:7C:89:45:B2:E8:7F:1A:3D:9C:5E:8B",
expiryDate: "2025-11-12"
},
{
type: "Certificate",
friendlyName: "web.company.com",
serialNumber: "8A:7F:1C:E4:92:B5:D3:68:F1:A2:7E:9B:4C:6D:5A:3F",
expiryDate: "2025-11-10"
}
],
alertName: "My PKI Alert",
alertName: "Production SSL Certificate Expiration Alert",
projectId: "c3b0ef29-915b-4cb1-8684-65b91b7fe02d",
siteUrl: "https://infisical.com"
} as PkiExpirationAlertTemplateProps;