Add PKI Alerting v2

This commit is contained in:
Carlos Monastyrski
2025-11-06 02:24:19 -03:00
parent 7afb76e542
commit 33617dbaa7
37 changed files with 5253 additions and 51 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";
@@ -353,6 +354,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
@@ -266,9 +266,21 @@ import {
TOrgRoles,
TOrgRolesInsert,
TOrgRolesUpdate,
TPkiAlertChannels,
TPkiAlertChannelsInsert,
TPkiAlertChannelsUpdate,
TPkiAlertHistory,
TPkiAlertHistoryCertificate,
TPkiAlertHistoryCertificateInsert,
TPkiAlertHistoryCertificateUpdate,
TPkiAlertHistoryInsert,
TPkiAlertHistoryUpdate,
TPkiAlerts,
TPkiAlertsInsert,
TPkiAlertsUpdate,
TPkiAlertsV2,
TPkiAlertsV2Insert,
TPkiAlertsV2Update,
TPkiApiEnrollmentConfigs,
TPkiApiEnrollmentConfigsInsert,
TPkiApiEnrollmentConfigsUpdate,
@@ -725,6 +737,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,89 @@
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("slug").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.index("eventType");
t.index("enabled");
t.unique(["slug", "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("notificationSent").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
@@ -92,7 +92,11 @@ export * from "./pam-accounts";
export * from "./pam-folders";
export * from "./pam-resources";
export * from "./pam-sessions";
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
@@ -29,6 +29,10 @@ export enum TableName {
PkiApiEnrollmentConfig = "pki_api_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(),
notificationSent: 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(),
slug: 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;
};
}
+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
@@ -267,6 +267,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";
@@ -547,6 +552,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);
@@ -1801,6 +1809,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
@@ -2355,6 +2378,7 @@ export const registerRoutes = async (
await dailyReminderQueueService.startSecretReminderMigrationJob();
await dailyExpiringPkiItemAlert.startSendingAlerts();
await pkiSubscriberQueue.startDailyAutoRenewalJob();
await pkiAlertV2Queue.startDailyAlertProcessing();
await certificateV3Queue.init();
await kmsService.startService(hsmStatus);
await microsoftTeamsService.start();
@@ -2486,7 +2510,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,433 @@
import { z } from "zod";
import { EventType } from "@app/ee/services/audit-log/audit-log-types";
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,
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: ["PKI Alerts"],
body: z.object({
projectId: z.string().uuid().describe("Project ID"),
...CreatePkiAlertV2Schema.shape
}),
response: {
200: z.object({
id: z.string().uuid(),
slug: z.string(),
description: z.string().nullable(),
eventType: z.nativeEnum(PkiAlertEventType),
alertBefore: z.string(),
filters: z.array(z.any()),
enabled: z.boolean(),
channels: z.array(
z.object({
id: z.string().uuid(),
channelType: z.string(),
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.slug,
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: ["PKI Alerts"],
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(),
slug: z.string(),
description: z.string().nullable(),
eventType: z.nativeEnum(PkiAlertEventType),
alertBefore: z.string(),
filters: z.array(z.any()),
enabled: z.boolean(),
channels: z.array(
z.object({
id: z.string().uuid(),
channelType: z.string(),
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: ["PKI Alerts"],
params: z.object({
alertId: z.string().uuid().describe("Alert ID")
}),
response: {
200: z.object({
id: z.string().uuid(),
slug: z.string(),
description: z.string().nullable(),
eventType: z.nativeEnum(PkiAlertEventType),
alertBefore: z.string(),
filters: z.array(z.any()),
enabled: z.boolean(),
channels: z.array(
z.object({
id: z.string().uuid(),
channelType: z.string(),
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: ["PKI Alerts"],
params: z.object({
alertId: z.string().uuid().describe("Alert ID")
}),
body: UpdatePkiAlertV2Schema,
response: {
200: z.object({
id: z.string().uuid(),
slug: z.string(),
description: z.string().nullable(),
eventType: z.nativeEnum(PkiAlertEventType),
alertBefore: z.string(),
filters: z.array(z.any()),
enabled: z.boolean(),
channels: z.array(
z.object({
id: z.string().uuid(),
channelType: z.string(),
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.slug,
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: ["PKI Alerts"],
params: z.object({
alertId: z.string().uuid().describe("Alert ID")
}),
response: {
200: z.object({
id: z.string().uuid(),
slug: z.string(),
description: z.string().nullable(),
eventType: z.nativeEnum(PkiAlertEventType),
alertBefore: z.string(),
filters: z.array(z.any()),
enabled: z.boolean(),
channels: z.array(
z.object({
id: z.string().uuid(),
channelType: z.string(),
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: ["PKI Alerts"],
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(),
limit: z.number(),
offset: 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: readLimit
},
onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]),
schema: {
description: "Preview certificates that would match the given filter rules",
tags: ["PKI Alerts"],
body: z.object({
projectId: z.string().uuid().describe("Project ID"),
filters: z.array(PkiFilterRuleSchema),
alertBefore: z
.string()
.regex(/^\d+[dwmy]$/)
.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(),
limit: z.number(),
offset: 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;
}
});
};
@@ -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?: {
notificationSent?: boolean;
notificationError?: string;
}
): Promise<TPkiAlertHistory> => {
try {
return await db.transaction(async (tx) => {
const historyRecords = await tx(TableName.PkiAlertHistory)
.insert({
alertId,
notificationSent: options?.notificationSent || 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.notificationSent", 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,525 @@
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 findByProjectId = async (
projectId: string,
filters?: {
search?: string;
eventType?: string;
enabled?: boolean;
limit?: number;
offset?: number;
},
tx?: Knex
): Promise<TAlertWithChannels[]> => {
try {
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}.slug`, `%${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 alerts = (await alertQuery) as TPkiAlertsV2[];
if (alerts.length === 0) {
return [];
}
const alertIds = alerts.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 result: TAlertWithChannels[] = alerts.map((alert) => ({
...alert,
channels: channelsByAlertId[alert.id] || []
}));
return result;
} catch (error) {
throw new DatabaseError({ error, name: "FindByProjectId" });
}
};
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}.slug`, `%${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)
.leftJoin(
`${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.slug 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)
.leftJoin(
`${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,
countByProjectId,
getDistinctProjectIds,
findMatchingCertificates
};
};
@@ -0,0 +1,340 @@
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 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 => {
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 condition = `${columnName}::text ILIKE ?`;
if (index === 0) {
void builder.whereRaw(condition, [`%"${String(v)}"%`]);
} else {
void builder.orWhereRaw(condition, [`%"${String(v)}"%`]);
}
});
});
}
return query.whereRaw(`${columnName}::text ILIKE ?`, [`%"${String(value)}"%`]);
case PkiFilterOperator.CONTAINS:
return applySanFilter(query, { ...filter, operator: PkiFilterOperator.MATCHES });
case PkiFilterOperator.STARTS_WITH:
return query.whereRaw(`${columnName}::text ILIKE ?`, [`%"${String(value)}%`]);
case PkiFilterOperator.ENDS_WITH:
return query.whereRaw(`${columnName}::text ILIKE ?`, [`%${String(value)}"%`]);
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);
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,244 @@
/* eslint-disable no-await-in-loop */
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 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;
slug: 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;
slug: 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.slug}) with ${certificateIds.length} certificates`
);
}
alertsProcessed += 1;
} catch (error) {
logger.error(error, `Failed to process alert ${typedAlert.id} (${typedAlert.slug})`);
}
}
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.`
);
};
queueService.start(QueueName.DailyPkiAlertV2Processing, async () => {
logger.info(`${QueueName.DailyPkiAlertV2Processing}: queue task started`);
try {
await processDailyAlerts();
logger.info(`${QueueName.DailyPkiAlertV2Processing}: queue task completed successfully`);
} catch (error) {
logger.error(error, `${QueueName.DailyPkiAlertV2Processing}: queue task failed`);
throw error;
}
});
const startDailyAlertProcessing = async () => {
await queueService.stopRepeatableJob(
QueueName.DailyPkiAlertV2Processing,
QueueJobs.DailyPkiAlertV2Processing,
{ pattern: "* * * * *", utc: true },
QueueName.DailyPkiAlertV2Processing
);
await queueService.queue(QueueName.DailyPkiAlertV2Processing, QueueJobs.DailyPkiAlertV2Processing, undefined, {
delay: 5000,
jobId: QueueName.DailyPkiAlertV2Processing,
repeat: { pattern: "* * * * *", utc: true }
});
logger.info("Daily PKI alert processing job scheduled");
};
const stopDailyAlertProcessing = async () => {
await queueService.stopRepeatableJob(
QueueName.DailyPkiAlertV2Processing,
QueueJobs.DailyPkiAlertV2Processing,
{ pattern: "* * * * *", utc: true },
QueueName.DailyPkiAlertV2Processing
);
logger.info("Daily PKI alert processing job stopped");
};
const triggerAlertProcessing = async () => {
await queueService.queue(QueueName.DailyPkiAlertV2Processing, QueueJobs.DailyPkiAlertV2Processing, undefined, {
delay: 1000
});
};
queueService.listen(QueueName.DailyPkiAlertV2Processing, "failed", (_, err) => {
logger.error(err, `${QueueName.DailyPkiAlertV2Processing}: Daily PKI alert processing failed`);
});
return {
startDailyAlertProcessing,
stopDailyAlertProcessing,
triggerAlertProcessing,
processDailyAlerts
};
};
@@ -0,0 +1,496 @@
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"
| "countByProjectId"
| "findMatchingCertificates"
>;
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;
slug: 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,
slug: alert.slug,
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,
slug,
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'" });
}
const alert = await pkiAlertV2DAL.create({
projectId,
slug,
description,
eventType,
alertBefore,
filters,
enabled
});
const channelInserts = channels.map((channel) => ({
alertId: alert.id,
channelType: channel.channelType,
config: channel.config,
enabled: channel.enabled
}));
await pkiAlertChannelDAL.insertMany(channelInserts);
const completeAlert = await pkiAlertV2DAL.findByIdWithChannels(alert.id);
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 Promise.all([
pkiAlertV2DAL.findByProjectId(projectId, filters),
pkiAlertV2DAL.countByProjectId(projectId, { search, eventType, enabled })
]);
return {
alerts: alerts.map((alert) => formatAlertResponse(alert as TAlertWithChannels)),
total
};
};
const updateAlert = async ({
alertId,
slug,
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: {
slug?: string;
description?: string;
eventType?: PkiAlertEventType;
alertBefore?: string;
filters?: TPkiFilterRule[];
enabled?: boolean;
} = {};
if (slug !== undefined) updateData.slug = slug;
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;
alert = await pkiAlertV2DAL.updateById(alertId, updateData);
if (channels) {
await pkiAlertChannelDAL.deleteByAlertId(alertId);
const channelInserts = channels.map((channel) => ({
alertId,
channelType: channel.channelType,
config: channel.config,
enabled: channel.enabled
}));
await pkiAlertChannelDAL.insertMany(channelInserts);
}
const completeAlert = await pkiAlertV2DAL.findByIdWithChannels(alertId);
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,
limit,
offset
};
};
const listCurrentMatchingCertificates = async ({
projectId,
filters,
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);
const options: {
limit: number;
offset: number;
showPreview?: boolean;
} = {
limit,
offset,
showPreview: true
};
const result = await pkiAlertV2DAL.findMatchingCertificates(projectId, filters, options);
return {
certificates: result.certificates,
total: result.total,
limit,
offset
};
};
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 notificationSent = 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 { slug: string }).slug;
const emailPromises = emailChannels.map((channel) => {
const config = channel.config as TEmailChannelConfig;
return smtpService.sendMail({
recipients: config.recipients,
subjectLine: `Infisical PKI 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);
notificationSent = 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, {
notificationSent,
notificationError
});
};
return {
createAlert,
getAlertById,
listAlerts,
updateAlert,
deleteAlert,
listMatchingCertificates,
listCurrentMatchingCertificates,
sendAlertNotifications
};
};
@@ -0,0 +1,198 @@
import RE2 from "re2";
import { z } from "zod";
import { TGenericPermission } from "@app/lib/types";
const createSecureSlugValidator = () => {
const slugRegex = new RE2("^[a-z0-9]+(?:-[a-z0-9]+)*$");
return (value: string) => slugRegex.test(value);
};
const createSecureAlertBeforeValidator = () => {
const alertBeforeRegex = new RE2("^\\d+[dwmy]$");
return (value: string) => 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({
slug: z
.string()
.min(1)
.max(255)
.refine(createSecureSlugValidator(), "Must be a valid slug (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;
slug: 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;
limit: number;
offset: 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,65 @@ 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 daysText = alertBeforeDays === 1 ? "1 day" : `${alertBeforeDays} days`;
const certificateText = items.length === 1 ? "certificate" : "certificates";
return (
<BaseEmailWrapper
title="Infisical CA/Certificate Expiration Notice"
preview="One or more of your Infisical certificates is about to expire."
title="Certificate Expiration Notice"
preview={`${items.length} ${certificateText} expiring in ${daysText}`}
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.
<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]">
You have{" "}
<strong>
{items.length} {certificateText}
</strong>{" "}
expiring in {daysText}.
</Text>
<Text className="text-[14px] leading-[24px] mb-[4px]">
<strong>Expiring Items:</strong>
<Text className="text-[14px]">
Alert: <strong>{alertName}</strong>
</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 +81,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;