diff --git a/backend/src/@types/fastify.d.ts b/backend/src/@types/fastify.d.ts index c25d8d4d1..95de82591 100644 --- a/backend/src/@types/fastify.d.ts +++ b/backend/src/@types/fastify.d.ts @@ -7,7 +7,7 @@ import { TAccessApprovalPolicyServiceFactory } from "@app/ee/services/access-app import { TAccessApprovalRequestServiceFactory } from "@app/ee/services/access-approval-request/access-approval-request-types"; import { TAssumePrivilegeServiceFactory } from "@app/ee/services/assume-privilege/assume-privilege-types"; import { TAuditLogServiceFactory, TCreateAuditLogDTO } from "@app/ee/services/audit-log/audit-log-types"; -import { TAuditLogStreamServiceFactory } from "@app/ee/services/audit-log-stream/audit-log-stream-types"; +import { TAuditLogStreamServiceFactory } from "@app/ee/services/audit-log-stream/audit-log-stream-service"; import { TCertificateAuthorityCrlServiceFactory } from "@app/ee/services/certificate-authority-crl/certificate-authority-crl-types"; import { TCertificateEstServiceFactory } from "@app/ee/services/certificate-est/certificate-est-service"; import { TDynamicSecretServiceFactory } from "@app/ee/services/dynamic-secret/dynamic-secret-types"; diff --git a/backend/src/db/migrations/20250903191434_audit-log-stream-v2.ts b/backend/src/db/migrations/20250903191434_audit-log-stream-v2.ts new file mode 100644 index 000000000..21d652c86 --- /dev/null +++ b/backend/src/db/migrations/20250903191434_audit-log-stream-v2.ts @@ -0,0 +1,210 @@ +import { Knex } from "knex"; + +import { inMemoryKeyStore } from "@app/keystore/memory"; +import { crypto } from "@app/lib/crypto/cryptography"; +import { KmsDataKey } from "@app/services/kms/kms-types"; +import { superAdminDALFactory } from "@app/services/super-admin/super-admin-dal"; + +import { SecretKeyEncoding, TableName } from "../schemas"; +import { getMigrationEnvConfig } from "./utils/env-config"; +import { createCircularCache } from "./utils/ring-buffer"; +import { getMigrationEncryptionServices } from "./utils/services"; + +const BATCH_SIZE = 500; +export async function up(knex: Knex): Promise { + if (await knex.schema.hasTable(TableName.AuditLogStream)) { + const hasProvider = await knex.schema.hasColumn(TableName.AuditLogStream, "provider"); + const hasEncryptedCredentials = await knex.schema.hasColumn(TableName.AuditLogStream, "encryptedCredentials"); + + await knex.schema.alterTable(TableName.AuditLogStream, (t) => { + if (!hasProvider) t.string("provider").notNullable().defaultTo("custom"); + if (!hasEncryptedCredentials) t.binary("encryptedCredentials"); + + // This column will no longer be used but we're not dropping it so that we can have a backup in case the migration goes wrong + t.string("url").nullable().alter(); + }); + + const superAdminDAL = superAdminDALFactory(knex); + const envConfig = await getMigrationEnvConfig(superAdminDAL); + const keyStore = inMemoryKeyStore(); + + const { kmsService } = await getMigrationEncryptionServices({ envConfig, keyStore, db: knex }); + + const orgEncryptionRingBuffer = + createCircularCache>>(25); + + const logStreams = await knex(TableName.AuditLogStream).select( + "id", + "orgId", + + "url", + "encryptedHeadersAlgorithm", + "encryptedHeadersCiphertext", + "encryptedHeadersIV", + "encryptedHeadersKeyEncoding", + "encryptedHeadersTag" + ); + + const updatedLogStreams = await Promise.all( + logStreams.map(async (el) => { + let orgKmsService = orgEncryptionRingBuffer.getItem(el.orgId); + if (!orgKmsService) { + orgKmsService = await kmsService.createCipherPairWithDataKey( + { + type: KmsDataKey.Organization, + orgId: el.orgId + }, + knex + ); + orgEncryptionRingBuffer.push(el.orgId, orgKmsService); + } + + const provider = "custom"; + let credentials; + + if ( + el.encryptedHeadersTag && + el.encryptedHeadersIV && + el.encryptedHeadersCiphertext && + el.encryptedHeadersKeyEncoding + ) { + const decryptedHeaders = crypto + .encryption() + .symmetric() + .decryptWithRootEncryptionKey({ + tag: el.encryptedHeadersTag, + iv: el.encryptedHeadersIV, + ciphertext: el.encryptedHeadersCiphertext, + keyEncoding: el.encryptedHeadersKeyEncoding as SecretKeyEncoding + }); + + credentials = { + url: el.url, + headers: JSON.parse(decryptedHeaders) + }; + } else { + credentials = { + url: el.url, + headers: [] + }; + } + + const encryptedCredentials = orgKmsService.encryptor({ + plainText: Buffer.from(JSON.stringify(credentials), "utf8") + }).cipherTextBlob; + + return { + id: el.id, + orgId: el.orgId, + url: el.url, + provider, + encryptedCredentials + }; + }) + ); + + for (let i = 0; i < updatedLogStreams.length; i += BATCH_SIZE) { + // eslint-disable-next-line no-await-in-loop + await knex(TableName.AuditLogStream) + .insert(updatedLogStreams.slice(i, i + BATCH_SIZE)) + .onConflict("id") + .merge(); + } + + await knex.schema.alterTable(TableName.AuditLogStream, (t) => { + t.binary("encryptedCredentials").notNullable().alter(); + }); + } +} + +// IMPORTANT: The down migration does not utilize the existing "url" and encrypted header columns +// because we're taking the latest data from the credentials column and re-encrypting it into relevant columns +// +// If this down migration was to fail, you can fall-back to the existing URL and encrypted header columns to retrieve +// data that was created prior to this migration + +export async function down(knex: Knex): Promise { + if (await knex.schema.hasTable(TableName.AuditLogStream)) { + const superAdminDAL = superAdminDALFactory(knex); + const envConfig = await getMigrationEnvConfig(superAdminDAL); + const keyStore = inMemoryKeyStore(); + + const { kmsService } = await getMigrationEncryptionServices({ envConfig, keyStore, db: knex }); + + const orgEncryptionRingBuffer = + createCircularCache>>(25); + + const logStreamsToRevert = await knex(TableName.AuditLogStream) + .select("id", "orgId", "encryptedCredentials") + .where("provider", "custom") + .whereNotNull("encryptedCredentials"); + + const updatedLogStreams = await Promise.all( + logStreamsToRevert.map(async (el) => { + let orgKmsService = orgEncryptionRingBuffer.getItem(el.orgId); + if (!orgKmsService) { + orgKmsService = await kmsService.createCipherPairWithDataKey( + { + type: KmsDataKey.Organization, + orgId: el.orgId + }, + knex + ); + orgEncryptionRingBuffer.push(el.orgId, orgKmsService); + } + + const decryptedCredentials = orgKmsService + .decryptor({ + cipherTextBlob: el.encryptedCredentials + }) + .toString(); + + const credentials: { url: string; headers: { key: string; value: string }[] } = + JSON.parse(decryptedCredentials); + + const originalUrl: string = credentials.url; + + const encryptedHeadersResult = crypto + .encryption() + .symmetric() + .encryptWithRootEncryptionKey(JSON.stringify(credentials.headers), envConfig); + + const encryptedHeadersAlgorithm: string = encryptedHeadersResult.algorithm; + const encryptedHeadersCiphertext: string = encryptedHeadersResult.ciphertext; + const encryptedHeadersIV: string = encryptedHeadersResult.iv; + const encryptedHeadersKeyEncoding: string = encryptedHeadersResult.encoding; + const encryptedHeadersTag: string = encryptedHeadersResult.tag; + + return { + id: el.id, + orgId: el.orgId, + encryptedCredentials: el.encryptedCredentials, + + url: originalUrl, + encryptedHeadersAlgorithm, + encryptedHeadersCiphertext, + encryptedHeadersIV, + encryptedHeadersKeyEncoding, + encryptedHeadersTag + }; + }) + ); + + for (let i = 0; i < updatedLogStreams.length; i += BATCH_SIZE) { + // eslint-disable-next-line no-await-in-loop + await knex(TableName.AuditLogStream) + .insert(updatedLogStreams.slice(i, i + BATCH_SIZE)) + .onConflict("id") + .merge(); + } + + await knex(TableName.AuditLogStream).whereNot("provider", "custom").orWhereNull("url").del(); + + await knex.schema.alterTable(TableName.AuditLogStream, (t) => { + t.string("url").notNullable().alter(); + + t.dropColumn("provider"); + t.dropColumn("encryptedCredentials"); + }); + } +} diff --git a/backend/src/db/schemas/audit-log-streams.ts b/backend/src/db/schemas/audit-log-streams.ts index 901dd8d27..a3f6bafba 100644 --- a/backend/src/db/schemas/audit-log-streams.ts +++ b/backend/src/db/schemas/audit-log-streams.ts @@ -5,11 +5,13 @@ import { z } from "zod"; +import { zodBuffer } from "@app/lib/zod"; + import { TImmutableDBKeys } from "./models"; export const AuditLogStreamsSchema = z.object({ id: z.string().uuid(), - url: z.string(), + url: z.string().nullable().optional(), encryptedHeadersCiphertext: z.string().nullable().optional(), encryptedHeadersIV: z.string().nullable().optional(), encryptedHeadersTag: z.string().nullable().optional(), @@ -17,7 +19,9 @@ export const AuditLogStreamsSchema = z.object({ encryptedHeadersKeyEncoding: z.string().nullable().optional(), orgId: z.string().uuid(), createdAt: z.date(), - updatedAt: z.date() + updatedAt: z.date(), + provider: z.string().default("custom"), + encryptedCredentials: zodBuffer }); export type TAuditLogStreams = z.infer; diff --git a/backend/src/ee/routes/v1/audit-log-stream-router.ts b/backend/src/ee/routes/v1/audit-log-stream-router.ts deleted file mode 100644 index 17bd9e64b..000000000 --- a/backend/src/ee/routes/v1/audit-log-stream-router.ts +++ /dev/null @@ -1,215 +0,0 @@ -import { z } from "zod"; - -import { AUDIT_LOG_STREAMS } from "@app/lib/api-docs"; -import { readLimit } from "@app/server/config/rateLimiter"; -import { verifyAuth } from "@app/server/plugins/auth/verify-auth"; -import { SanitizedAuditLogStreamSchema } from "@app/server/routes/sanitizedSchemas"; -import { AuthMode } from "@app/services/auth/auth-type"; - -export const registerAuditLogStreamRouter = async (server: FastifyZodProvider) => { - server.route({ - method: "POST", - url: "/", - config: { - rateLimit: readLimit - }, - schema: { - description: "Create an Audit Log Stream.", - security: [ - { - bearerAuth: [] - } - ], - body: z.object({ - url: z.string().min(1).describe(AUDIT_LOG_STREAMS.CREATE.url), - headers: z - .object({ - key: z.string().min(1).trim().describe(AUDIT_LOG_STREAMS.CREATE.headers.key), - value: z.string().min(1).trim().describe(AUDIT_LOG_STREAMS.CREATE.headers.value) - }) - .describe(AUDIT_LOG_STREAMS.CREATE.headers.desc) - .array() - .optional() - }), - response: { - 200: z.object({ - auditLogStream: SanitizedAuditLogStreamSchema - }) - } - }, - onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]), - handler: async (req) => { - const auditLogStream = await server.services.auditLogStream.create({ - actorId: req.permission.id, - actor: req.permission.type, - actorOrgId: req.permission.orgId, - actorAuthMethod: req.permission.authMethod, - url: req.body.url, - headers: req.body.headers - }); - - return { auditLogStream }; - } - }); - - server.route({ - method: "PATCH", - url: "/:id", - config: { - rateLimit: readLimit - }, - schema: { - description: "Update an Audit Log Stream by ID.", - security: [ - { - bearerAuth: [] - } - ], - params: z.object({ - id: z.string().describe(AUDIT_LOG_STREAMS.UPDATE.id) - }), - body: z.object({ - url: z.string().optional().describe(AUDIT_LOG_STREAMS.UPDATE.url), - headers: z - .object({ - key: z.string().min(1).trim().describe(AUDIT_LOG_STREAMS.UPDATE.headers.key), - value: z.string().min(1).trim().describe(AUDIT_LOG_STREAMS.UPDATE.headers.value) - }) - .describe(AUDIT_LOG_STREAMS.UPDATE.headers.desc) - .array() - .optional() - }), - response: { - 200: z.object({ - auditLogStream: SanitizedAuditLogStreamSchema - }) - } - }, - onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]), - handler: async (req) => { - const auditLogStream = await server.services.auditLogStream.updateById({ - actorId: req.permission.id, - actor: req.permission.type, - actorOrgId: req.permission.orgId, - actorAuthMethod: req.permission.authMethod, - id: req.params.id, - url: req.body.url, - headers: req.body.headers - }); - - return { auditLogStream }; - } - }); - - server.route({ - method: "DELETE", - url: "/:id", - config: { - rateLimit: readLimit - }, - schema: { - description: "Delete an Audit Log Stream by ID.", - security: [ - { - bearerAuth: [] - } - ], - params: z.object({ - id: z.string().describe(AUDIT_LOG_STREAMS.DELETE.id) - }), - response: { - 200: z.object({ - auditLogStream: SanitizedAuditLogStreamSchema - }) - } - }, - onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]), - handler: async (req) => { - const auditLogStream = await server.services.auditLogStream.deleteById({ - actorId: req.permission.id, - actor: req.permission.type, - actorOrgId: req.permission.orgId, - actorAuthMethod: req.permission.authMethod, - id: req.params.id - }); - - return { auditLogStream }; - } - }); - - server.route({ - method: "GET", - url: "/:id", - config: { - rateLimit: readLimit - }, - schema: { - description: "Get an Audit Log Stream by ID.", - security: [ - { - bearerAuth: [] - } - ], - params: z.object({ - id: z.string().describe(AUDIT_LOG_STREAMS.GET_BY_ID.id) - }), - response: { - 200: z.object({ - auditLogStream: SanitizedAuditLogStreamSchema.extend({ - headers: z - .object({ - key: z.string(), - value: z.string() - }) - .array() - .optional() - }) - }) - } - }, - onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]), - handler: async (req) => { - const auditLogStream = await server.services.auditLogStream.getById({ - actorId: req.permission.id, - actor: req.permission.type, - actorOrgId: req.permission.orgId, - actorAuthMethod: req.permission.authMethod, - id: req.params.id - }); - - return { auditLogStream }; - } - }); - - server.route({ - method: "GET", - url: "/", - config: { - rateLimit: readLimit - }, - schema: { - description: "List Audit Log Streams.", - security: [ - { - bearerAuth: [] - } - ], - response: { - 200: z.object({ - auditLogStreams: SanitizedAuditLogStreamSchema.array() - }) - } - }, - onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]), - handler: async (req) => { - const auditLogStreams = await server.services.auditLogStream.list({ - actorId: req.permission.id, - actor: req.permission.type, - actorOrgId: req.permission.orgId, - actorAuthMethod: req.permission.authMethod - }); - - return { auditLogStreams }; - } - }); -}; diff --git a/backend/src/ee/routes/v1/audit-log-stream-routers/audit-log-stream-endpoints.ts b/backend/src/ee/routes/v1/audit-log-stream-routers/audit-log-stream-endpoints.ts new file mode 100644 index 000000000..816fc62fb --- /dev/null +++ b/backend/src/ee/routes/v1/audit-log-stream-routers/audit-log-stream-endpoints.ts @@ -0,0 +1,142 @@ +import { z } from "zod"; + +import { LogProvider } from "@app/ee/services/audit-log-stream/audit-log-stream-enums"; +import { TAuditLogStream } from "@app/ee/services/audit-log-stream/audit-log-stream-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"; + +export const registerAuditLogStreamEndpoints = ({ + server, + provider, + createSchema, + updateSchema, + sanitizedResponseSchema +}: { + server: FastifyZodProvider; + provider: LogProvider; + createSchema: z.ZodType<{ + credentials: T["credentials"]; + }>; + updateSchema: z.ZodType<{ + credentials: T["credentials"]; + }>; + sanitizedResponseSchema: z.ZodTypeAny; +}) => { + server.route({ + method: "GET", + url: "/:logStreamId", + config: { + rateLimit: readLimit + }, + schema: { + params: z.object({ + logStreamId: z.string().uuid() + }), + response: { + 200: z.object({ + auditLogStream: sanitizedResponseSchema + }) + } + }, + onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]), + handler: async (req) => { + const { logStreamId } = req.params; + + const auditLogStream = await server.services.auditLogStream.getById(logStreamId, provider, req.permission); + + return { auditLogStream }; + } + }); + + server.route({ + method: "POST", + url: "/", + config: { + rateLimit: writeLimit + }, + schema: { + body: createSchema, + response: { + 200: z.object({ + auditLogStream: sanitizedResponseSchema + }) + } + }, + onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]), + handler: async (req) => { + const { credentials } = req.body; + + const auditLogStream = await server.services.auditLogStream.create( + { + provider, + credentials + }, + req.permission + ); + + return { auditLogStream }; + } + }); + + server.route({ + method: "PATCH", + url: "/:logStreamId", + config: { + rateLimit: writeLimit + }, + schema: { + params: z.object({ + logStreamId: z.string().uuid() + }), + body: updateSchema, + response: { + 200: z.object({ + auditLogStream: sanitizedResponseSchema + }) + } + }, + onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]), + handler: async (req) => { + const { logStreamId } = req.params; + const { credentials } = req.body; + + const auditLogStream = await server.services.auditLogStream.updateById( + { + logStreamId, + provider, + credentials + }, + req.permission + ); + + return { auditLogStream }; + } + }); + + server.route({ + method: "DELETE", + url: "/:logStreamId", + config: { + rateLimit: writeLimit + }, + schema: { + params: z.object({ + logStreamId: z.string().uuid() + }), + response: { + 200: z.object({ + auditLogStream: sanitizedResponseSchema + }) + } + }, + onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]), + handler: async (req) => { + const { logStreamId } = req.params; + + const auditLogStream = await server.services.auditLogStream.deleteById(logStreamId, provider, req.permission); + + return { auditLogStream }; + } + }); +}; diff --git a/backend/src/ee/routes/v1/audit-log-stream-routers/audit-log-stream-router.ts b/backend/src/ee/routes/v1/audit-log-stream-routers/audit-log-stream-router.ts new file mode 100644 index 000000000..bf10e48bf --- /dev/null +++ b/backend/src/ee/routes/v1/audit-log-stream-routers/audit-log-stream-router.ts @@ -0,0 +1,73 @@ +import { z } from "zod"; + +import { + CustomProviderListItemSchema, + SanitizedCustomProviderSchema +} from "@app/ee/services/audit-log-stream/custom/custom-provider-schemas"; +import { + DatadogProviderListItemSchema, + SanitizedDatadogProviderSchema +} from "@app/ee/services/audit-log-stream/datadog/datadog-provider-schemas"; +import { + SanitizedSplunkProviderSchema, + SplunkProviderListItemSchema +} from "@app/ee/services/audit-log-stream/splunk/splunk-provider-schemas"; +import { readLimit } from "@app/server/config/rateLimiter"; +import { verifyAuth } from "@app/server/plugins/auth/verify-auth"; +import { AuthMode } from "@app/services/auth/auth-type"; + +const SanitizedAuditLogStreamSchema = z.union([ + SanitizedCustomProviderSchema, + SanitizedDatadogProviderSchema, + SanitizedSplunkProviderSchema +]); + +const ProviderOptionsSchema = z.discriminatedUnion("provider", [ + CustomProviderListItemSchema, + DatadogProviderListItemSchema, + SplunkProviderListItemSchema +]); + +export const registerAuditLogStreamRouter = async (server: FastifyZodProvider) => { + server.route({ + method: "GET", + url: "/options", + config: { + rateLimit: readLimit + }, + schema: { + response: { + 200: z.object({ + providerOptions: ProviderOptionsSchema.array() + }) + } + }, + onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]), + handler: () => { + const providerOptions = server.services.auditLogStream.listProviderOptions(); + + return { providerOptions }; + } + }); + + server.route({ + method: "GET", + url: "/", + config: { + rateLimit: readLimit + }, + schema: { + response: { + 200: z.object({ + auditLogStreams: SanitizedAuditLogStreamSchema.array() + }) + } + }, + onRequest: verifyAuth([AuthMode.JWT, AuthMode.IDENTITY_ACCESS_TOKEN]), + handler: async (req) => { + const auditLogStreams = await server.services.auditLogStream.list(req.permission); + + return { auditLogStreams }; + } + }); +}; diff --git a/backend/src/ee/routes/v1/audit-log-stream-routers/index.ts b/backend/src/ee/routes/v1/audit-log-stream-routers/index.ts new file mode 100644 index 000000000..2f19e7522 --- /dev/null +++ b/backend/src/ee/routes/v1/audit-log-stream-routers/index.ts @@ -0,0 +1,48 @@ +import { LogProvider } from "@app/ee/services/audit-log-stream/audit-log-stream-enums"; +import { + CustomProviderSchema, + SanitizedCustomProviderSchema +} from "@app/ee/services/audit-log-stream/custom/custom-provider-schemas"; +import { + DatadogProviderSchema, + SanitizedDatadogProviderSchema +} from "@app/ee/services/audit-log-stream/datadog/datadog-provider-schemas"; +import { + SanitizedSplunkProviderSchema, + SplunkProviderSchema +} from "@app/ee/services/audit-log-stream/splunk/splunk-provider-schemas"; + +import { registerAuditLogStreamEndpoints } from "./audit-log-stream-endpoints"; + +export * from "./audit-log-stream-router"; + +export const AUDIT_LOG_STREAM_REGISTER_ROUTER_MAP: Record Promise> = + { + [LogProvider.Custom]: async (server: FastifyZodProvider) => { + registerAuditLogStreamEndpoints({ + server, + provider: LogProvider.Custom, + sanitizedResponseSchema: SanitizedCustomProviderSchema, + createSchema: CustomProviderSchema, + updateSchema: CustomProviderSchema + }); + }, + [LogProvider.Datadog]: async (server: FastifyZodProvider) => { + registerAuditLogStreamEndpoints({ + server, + provider: LogProvider.Datadog, + sanitizedResponseSchema: SanitizedDatadogProviderSchema, + createSchema: DatadogProviderSchema, + updateSchema: DatadogProviderSchema + }); + }, + [LogProvider.Splunk]: async (server: FastifyZodProvider) => { + registerAuditLogStreamEndpoints({ + server, + provider: LogProvider.Splunk, + sanitizedResponseSchema: SanitizedSplunkProviderSchema, + createSchema: SplunkProviderSchema, + updateSchema: SplunkProviderSchema + }); + } + }; diff --git a/backend/src/ee/routes/v1/index.ts b/backend/src/ee/routes/v1/index.ts index ab9503f58..35179c845 100644 --- a/backend/src/ee/routes/v1/index.ts +++ b/backend/src/ee/routes/v1/index.ts @@ -3,7 +3,7 @@ import { registerProjectTemplateRouter } from "@app/ee/routes/v1/project-templat import { registerAccessApprovalPolicyRouter } from "./access-approval-policy-router"; import { registerAccessApprovalRequestRouter } from "./access-approval-request-router"; import { registerAssumePrivilegeRouter } from "./assume-privilege-router"; -import { registerAuditLogStreamRouter } from "./audit-log-stream-router"; +import { AUDIT_LOG_STREAM_REGISTER_ROUTER_MAP, registerAuditLogStreamRouter } from "./audit-log-stream-routers"; import { registerCaCrlRouter } from "./certificate-authority-crl-router"; import { registerDynamicSecretLeaseRouter } from "./dynamic-secret-lease-router"; import { registerKubernetesDynamicSecretLeaseRouter } from "./dynamic-secret-lease-routers/kubernetes-lease-router"; @@ -114,7 +114,19 @@ export const registerV1EERoutes = async (server: FastifyZodProvider) => { await server.register(registerSecretRouter, { prefix: "/secrets" }); await server.register(registerSecretVersionRouter, { prefix: "/secret" }); await server.register(registerGroupRouter, { prefix: "/groups" }); - await server.register(registerAuditLogStreamRouter, { prefix: "/audit-log-streams" }); + + await server.register( + async (auditLogStreamRouter) => { + await auditLogStreamRouter.register(registerAuditLogStreamRouter); + + // Provider-specific endpoints + for await (const [provider, router] of Object.entries(AUDIT_LOG_STREAM_REGISTER_ROUTER_MAP)) { + await auditLogStreamRouter.register(router, { prefix: `/${provider}` }); + } + }, + { prefix: "/audit-log-streams" } + ); + await server.register(registerUserAdditionalPrivilegeRouter, { prefix: "/user-project-additional-privilege" }); await server.register( async (privilegeRouter) => { diff --git a/backend/src/ee/services/audit-log-stream/audit-log-stream-enums.ts b/backend/src/ee/services/audit-log-stream/audit-log-stream-enums.ts new file mode 100644 index 000000000..9c03dc2bd --- /dev/null +++ b/backend/src/ee/services/audit-log-stream/audit-log-stream-enums.ts @@ -0,0 +1,5 @@ +export enum LogProvider { + Datadog = "datadog", + Splunk = "splunk", + Custom = "custom" +} diff --git a/backend/src/ee/services/audit-log-stream/audit-log-stream-factory.ts b/backend/src/ee/services/audit-log-stream/audit-log-stream-factory.ts new file mode 100644 index 000000000..f70dc02a0 --- /dev/null +++ b/backend/src/ee/services/audit-log-stream/audit-log-stream-factory.ts @@ -0,0 +1,13 @@ +import { LogProvider } from "./audit-log-stream-enums"; +import { TAuditLogStreamCredentials, TLogStreamFactory } from "./audit-log-stream-types"; +import { CustomProviderFactory } from "./custom/custom-provider-factory"; +import { DatadogProviderFactory } from "./datadog/datadog-provider-factory"; +import { SplunkProviderFactory } from "./splunk/splunk-provider-factory"; + +type TLogStreamFactoryImplementation = TLogStreamFactory; + +export const LOG_STREAM_FACTORY_MAP: Record = { + [LogProvider.Datadog]: DatadogProviderFactory as TLogStreamFactoryImplementation, + [LogProvider.Splunk]: SplunkProviderFactory as TLogStreamFactoryImplementation, + [LogProvider.Custom]: CustomProviderFactory as TLogStreamFactoryImplementation +}; diff --git a/backend/src/ee/services/audit-log-stream/audit-log-stream-fns.ts b/backend/src/ee/services/audit-log-stream/audit-log-stream-fns.ts index dc93f238e..bf6e4d47d 100644 --- a/backend/src/ee/services/audit-log-stream/audit-log-stream-fns.ts +++ b/backend/src/ee/services/audit-log-stream/audit-log-stream-fns.ts @@ -1,21 +1,70 @@ -export function providerSpecificPayload(url: string) { - const { hostname } = new URL(url); +import { TAuditLogStreams } from "@app/db/schemas"; +import { TKmsServiceFactory } from "@app/services/kms/kms-service"; +import { KmsDataKey } from "@app/services/kms/kms-types"; - const payload: Record = {}; +import { TAuditLogStream, TAuditLogStreamCredentials } from "./audit-log-stream-types"; +import { getCustomProviderListItem } from "./custom/custom-provider-fns"; +import { getDatadogProviderListItem } from "./datadog/datadog-provider-fns"; +import { getSplunkProviderListItem } from "./splunk/splunk-provider-fns"; - switch (hostname) { - case "http-intake.logs.datadoghq.com": - case "http-intake.logs.us3.datadoghq.com": - case "http-intake.logs.us5.datadoghq.com": - case "http-intake.logs.datadoghq.eu": - case "http-intake.logs.ap1.datadoghq.com": - case "http-intake.logs.ddog-gov.com": - payload.ddsource = "infisical"; - payload.service = "audit-logs"; - break; - default: - break; - } +export const listProviderOptions = () => { + return [getCustomProviderListItem(), getDatadogProviderListItem(), getSplunkProviderListItem()].sort((a, b) => + a.name.localeCompare(b.name) + ); +}; - return payload; -} +export const encryptLogStreamCredentials = async ({ + orgId, + credentials, + kmsService +}: { + orgId: string; + credentials: TAuditLogStreamCredentials; + kmsService: Pick; +}) => { + const { encryptor } = await kmsService.createCipherPairWithDataKey({ + type: KmsDataKey.Organization, + orgId + }); + + const { cipherTextBlob: encryptedCredentialsBlob } = encryptor({ + plainText: Buffer.from(JSON.stringify(credentials)) + }); + + return encryptedCredentialsBlob; +}; + +export const decryptLogStreamCredentials = async ({ + orgId, + encryptedCredentials, + kmsService +}: { + orgId: string; + encryptedCredentials: Buffer; + kmsService: Pick; +}) => { + const { decryptor } = await kmsService.createCipherPairWithDataKey({ + type: KmsDataKey.Organization, + orgId + }); + + const decryptedPlainTextBlob = decryptor({ + cipherTextBlob: encryptedCredentials + }); + + return JSON.parse(decryptedPlainTextBlob.toString()) as TAuditLogStreamCredentials; +}; + +export const decryptLogStream = async ( + logStream: TAuditLogStreams, + kmsService: Pick +) => { + return { + ...logStream, + credentials: await decryptLogStreamCredentials({ + encryptedCredentials: logStream.encryptedCredentials, + orgId: logStream.orgId, + kmsService + }) + } as TAuditLogStream; +}; diff --git a/backend/src/ee/services/audit-log-stream/audit-log-stream-service.ts b/backend/src/ee/services/audit-log-stream/audit-log-stream-service.ts index 46d2782b3..773f32407 100644 --- a/backend/src/ee/services/audit-log-stream/audit-log-stream-service.ts +++ b/backend/src/ee/services/audit-log-stream/audit-log-stream-service.ts @@ -1,242 +1,215 @@ import { ForbiddenError } from "@casl/ability"; -import { RawAxiosRequestHeaders } from "axios"; +import { AxiosError } from "axios"; -import { SecretKeyEncoding } from "@app/db/schemas"; -import { getConfig } from "@app/lib/config/env"; -import { request } from "@app/lib/config/request"; -import { crypto } from "@app/lib/crypto/cryptography"; -import { BadRequestError, NotFoundError, UnauthorizedError } from "@app/lib/errors"; -import { blockLocalAndPrivateIpAddresses } from "@app/lib/validator"; +import { TAuditLogs } from "@app/db/schemas"; +import { + decryptLogStream, + decryptLogStreamCredentials, + encryptLogStreamCredentials, + listProviderOptions +} from "@app/ee/services/audit-log-stream/audit-log-stream-fns"; +import { BadRequestError, NotFoundError } from "@app/lib/errors"; +import { logger } from "@app/lib/logger"; +import { OrgServiceActor } from "@app/lib/types"; +import { TKmsServiceFactory } from "@app/services/kms/kms-service"; -import { AUDIT_LOG_STREAM_TIMEOUT } from "../audit-log/audit-log-queue"; import { TLicenseServiceFactory } from "../license/license-service"; import { OrgPermissionActions, OrgPermissionSubjects } from "../permission/org-permission"; import { TPermissionServiceFactory } from "../permission/permission-service-types"; import { TAuditLogStreamDALFactory } from "./audit-log-stream-dal"; -import { providerSpecificPayload } from "./audit-log-stream-fns"; -import { LogStreamHeaders, TAuditLogStreamServiceFactory } from "./audit-log-stream-types"; +import { LogProvider } from "./audit-log-stream-enums"; +import { LOG_STREAM_FACTORY_MAP } from "./audit-log-stream-factory"; +import { TAuditLogStream, TCreateAuditLogStreamDTO, TUpdateAuditLogStreamDTO } from "./audit-log-stream-types"; -type TAuditLogStreamServiceFactoryDep = { +export type TAuditLogStreamServiceFactoryDep = { auditLogStreamDAL: TAuditLogStreamDALFactory; permissionService: Pick; licenseService: Pick; + kmsService: Pick; }; +export type TAuditLogStreamServiceFactory = ReturnType; + export const auditLogStreamServiceFactory = ({ auditLogStreamDAL, permissionService, - licenseService -}: TAuditLogStreamServiceFactoryDep): TAuditLogStreamServiceFactory => { - const create: TAuditLogStreamServiceFactory["create"] = async ({ - url, - actor, - headers = [], - actorId, - actorOrgId, - actorAuthMethod - }) => { - if (!actorOrgId) throw new UnauthorizedError({ message: "No organization ID attached to authentication token" }); - - const plan = await licenseService.getPlan(actorOrgId); + licenseService, + kmsService +}: TAuditLogStreamServiceFactoryDep) => { + const create = async ({ provider, credentials }: TCreateAuditLogStreamDTO, actor: OrgServiceActor) => { + const plan = await licenseService.getPlan(actor.orgId); if (!plan.auditLogStreams) { throw new BadRequestError({ - message: "Failed to create audit log streams due to plan restriction. Upgrade plan to create group." + message: "Failed to create Audit Log Stream: Plan restriction. Upgrade plan to continue." }); } const { permission } = await permissionService.getOrgPermission( - actor, - actorId, - actorOrgId, - actorAuthMethod, - actorOrgId + actor.type, + actor.id, + actor.orgId, + actor.authMethod, + actor.orgId ); + ForbiddenError.from(permission).throwUnlessCan(OrgPermissionActions.Create, OrgPermissionSubjects.Settings); - const appCfg = getConfig(); - if (appCfg.isCloud) await blockLocalAndPrivateIpAddresses(url); - - const totalStreams = await auditLogStreamDAL.find({ orgId: actorOrgId }); + const totalStreams = await auditLogStreamDAL.find({ orgId: actor.orgId }); if (totalStreams.length >= plan.auditLogStreamLimit) { throw new BadRequestError({ - message: - "Failed to create audit log streams due to plan limit reached. Kindly contact Infisical to add more streams." + message: "Failed to create Audit Log Stream: Plan limit reached. Contact Infisical to increase quota." }); } - // testing connection first - const streamHeaders: RawAxiosRequestHeaders = { "Content-Type": "application/json" }; - if (headers.length) - headers.forEach(({ key, value }) => { - streamHeaders[key] = value; - }); + const factory = LOG_STREAM_FACTORY_MAP[provider](); + const validatedCredentials = await factory.validateCredentials({ credentials }); - await request - .post( - url, - { ...providerSpecificPayload(url), ping: "ok" }, - { - headers: streamHeaders, - // request timeout - timeout: AUDIT_LOG_STREAM_TIMEOUT, - // connection timeout - signal: AbortSignal.timeout(AUDIT_LOG_STREAM_TIMEOUT) - } - ) - .catch((err) => { - throw new BadRequestError({ message: `Failed to connect with upstream source: ${(err as Error)?.message}` }); - }); + const encryptedCredentials = await encryptLogStreamCredentials({ + credentials: validatedCredentials, + orgId: actor.orgId, + kmsService + }); - const encryptedHeaders = headers - ? crypto.encryption().symmetric().encryptWithRootEncryptionKey(JSON.stringify(headers)) - : undefined; const logStream = await auditLogStreamDAL.create({ - orgId: actorOrgId, - url, - ...(encryptedHeaders - ? { - encryptedHeadersCiphertext: encryptedHeaders.ciphertext, - encryptedHeadersIV: encryptedHeaders.iv, - encryptedHeadersTag: encryptedHeaders.tag, - encryptedHeadersAlgorithm: encryptedHeaders.algorithm, - encryptedHeadersKeyEncoding: encryptedHeaders.encoding - } - : {}) + orgId: actor.orgId, + provider, + encryptedCredentials }); - return logStream; + + return { ...logStream, credentials: validatedCredentials } as TAuditLogStream; }; - const updateById: TAuditLogStreamServiceFactory["updateById"] = async ({ - id, - url, - actor, - headers = [], - actorId, - actorOrgId, - actorAuthMethod - }) => { - if (!actorOrgId) throw new UnauthorizedError({ message: "No organization ID attached to authentication token" }); - - const plan = await licenseService.getPlan(actorOrgId); - if (!plan.auditLogStreams) + const updateById = async ( + { logStreamId, provider, credentials }: TUpdateAuditLogStreamDTO, + actor: OrgServiceActor + ) => { + const plan = await licenseService.getPlan(actor.orgId); + if (!plan.auditLogStreams) { throw new BadRequestError({ - message: "Failed to update audit log streams due to plan restriction. Upgrade plan to create group." + message: "Failed to update Audit Log Stream: Plan restriction. Upgrade plan to continue." }); + } - const logStream = await auditLogStreamDAL.findById(id); - if (!logStream) throw new NotFoundError({ message: `Audit log stream with ID '${id}' not found` }); + const logStream = await auditLogStreamDAL.findById(logStreamId); + if (!logStream) throw new NotFoundError({ message: `Audit Log Stream with ID '${logStreamId}' not found` }); + + const { permission } = await permissionService.getOrgPermission( + actor.type, + actor.id, + actor.orgId, + actor.authMethod, + logStream.orgId + ); - const { orgId } = logStream; - const { permission } = await permissionService.getOrgPermission(actor, actorId, orgId, actorAuthMethod, actorOrgId); ForbiddenError.from(permission).throwUnlessCan(OrgPermissionActions.Edit, OrgPermissionSubjects.Settings); - const appCfg = getConfig(); - if (url && appCfg.isCloud) await blockLocalAndPrivateIpAddresses(url); - // testing connection first - const streamHeaders: RawAxiosRequestHeaders = { "Content-Type": "application/json" }; - if (headers.length) - headers.forEach(({ key, value }) => { - streamHeaders[key] = value; - }); + const factory = LOG_STREAM_FACTORY_MAP[provider](); + const validatedCredentials = await factory.validateCredentials({ credentials }); - await request - .post( - url || logStream.url, - { ...providerSpecificPayload(url || logStream.url), ping: "ok" }, - { - headers: streamHeaders, - // request timeout - timeout: AUDIT_LOG_STREAM_TIMEOUT, - // connection timeout - signal: AbortSignal.timeout(AUDIT_LOG_STREAM_TIMEOUT) - } - ) - .catch((err) => { - throw new Error(`Failed to connect with the source ${(err as Error)?.message}`); - }); - - const encryptedHeaders = headers - ? crypto.encryption().symmetric().encryptWithRootEncryptionKey(JSON.stringify(headers)) - : undefined; - const updatedLogStream = await auditLogStreamDAL.updateById(id, { - url, - ...(encryptedHeaders - ? { - encryptedHeadersCiphertext: encryptedHeaders.ciphertext, - encryptedHeadersIV: encryptedHeaders.iv, - encryptedHeadersTag: encryptedHeaders.tag, - encryptedHeadersAlgorithm: encryptedHeaders.algorithm, - encryptedHeadersKeyEncoding: encryptedHeaders.encoding - } - : {}) + const encryptedCredentials = await encryptLogStreamCredentials({ + credentials: validatedCredentials, + orgId: actor.orgId, + kmsService }); - return updatedLogStream; + + const updatedLogStream = await auditLogStreamDAL.updateById(logStreamId, { + encryptedCredentials + }); + + return { ...updatedLogStream, credentials: validatedCredentials } as TAuditLogStream; }; - const deleteById: TAuditLogStreamServiceFactory["deleteById"] = async ({ - id, - actor, - actorId, - actorOrgId, - actorAuthMethod - }) => { - if (!actorOrgId) throw new UnauthorizedError({ message: "No organization ID attached to authentication token" }); + const deleteById = async (logStreamId: string, provider: LogProvider, actor: OrgServiceActor) => { + const logStream = await auditLogStreamDAL.findById(logStreamId); + if (!logStream) throw new NotFoundError({ message: `Audit Log Stream with ID '${logStreamId}' not found` }); - const logStream = await auditLogStreamDAL.findById(id); - if (!logStream) throw new NotFoundError({ message: `Audit log stream with ID '${id}' not found` }); + const { permission } = await permissionService.getOrgPermission( + actor.type, + actor.id, + actor.orgId, + actor.authMethod, + logStream.orgId + ); - const { orgId } = logStream; - const { permission } = await permissionService.getOrgPermission(actor, actorId, orgId, actorAuthMethod, actorOrgId); ForbiddenError.from(permission).throwUnlessCan(OrgPermissionActions.Delete, OrgPermissionSubjects.Settings); - const deletedLogStream = await auditLogStreamDAL.deleteById(id); - return deletedLogStream; + if (logStream.provider !== provider) { + throw new BadRequestError({ + message: `Audit Log Stream with ID '${logStreamId}' is not for provider '${provider}'` + }); + } + + const deletedLogStream = await auditLogStreamDAL.deleteById(logStreamId); + + return decryptLogStream(deletedLogStream, kmsService); }; - const getById: TAuditLogStreamServiceFactory["getById"] = async ({ - id, - actor, - actorId, - actorOrgId, - actorAuthMethod - }) => { - const logStream = await auditLogStreamDAL.findById(id); - if (!logStream) throw new NotFoundError({ message: `Audit log stream with ID '${id}' not found` }); + const getById = async (logStreamId: string, provider: LogProvider, actor: OrgServiceActor) => { + const logStream = await auditLogStreamDAL.findById(logStreamId); - const { orgId } = logStream; - const { permission } = await permissionService.getOrgPermission(actor, actorId, orgId, actorAuthMethod, actorOrgId); - ForbiddenError.from(permission).throwUnlessCan(OrgPermissionActions.Read, OrgPermissionSubjects.Settings); + if (!logStream) throw new NotFoundError({ message: `Audit log stream with ID '${logStreamId}' not found` }); - const headers = - logStream?.encryptedHeadersCiphertext && logStream?.encryptedHeadersIV && logStream?.encryptedHeadersTag - ? (JSON.parse( - crypto - .encryption() - .symmetric() - .decryptWithRootEncryptionKey({ - tag: logStream.encryptedHeadersTag, - iv: logStream.encryptedHeadersIV, - ciphertext: logStream.encryptedHeadersCiphertext, - keyEncoding: logStream.encryptedHeadersKeyEncoding as SecretKeyEncoding - }) - ) as LogStreamHeaders[]) - : undefined; - - return { ...logStream, headers }; - }; - - const list: TAuditLogStreamServiceFactory["list"] = async ({ actor, actorId, actorOrgId, actorAuthMethod }) => { const { permission } = await permissionService.getOrgPermission( - actor, - actorId, - actorOrgId, - actorAuthMethod, - actorOrgId + actor.type, + actor.id, + logStream.orgId, + actor.authMethod, + actor.orgId ); + ForbiddenError.from(permission).throwUnlessCan(OrgPermissionActions.Read, OrgPermissionSubjects.Settings); - const logStreams = await auditLogStreamDAL.find({ orgId: actorOrgId }); - return logStreams; + if (logStream.provider !== provider) { + throw new BadRequestError({ + message: `Audit Log Stream with ID '${logStreamId}' is not for provider '${provider}'` + }); + } + + return decryptLogStream(logStream, kmsService); + }; + + const list = async (actor: OrgServiceActor) => { + const { permission } = await permissionService.getOrgPermission( + actor.type, + actor.id, + actor.orgId, + actor.authMethod, + actor.orgId + ); + + ForbiddenError.from(permission).throwUnlessCan(OrgPermissionActions.Read, OrgPermissionSubjects.Settings); + + const logStreams = await auditLogStreamDAL.find({ orgId: actor.orgId }); + + return Promise.all(logStreams.map((stream) => decryptLogStream(stream, kmsService))); + }; + + const streamLog = async (orgId: string, auditLog: TAuditLogs) => { + const logStreams = await auditLogStreamDAL.find({ orgId }); + await Promise.allSettled( + logStreams.map(async ({ provider, encryptedCredentials }) => { + const credentials = await decryptLogStreamCredentials({ + encryptedCredentials, + orgId, + kmsService + }); + + const factory = LOG_STREAM_FACTORY_MAP[provider as LogProvider](); + + try { + await factory.streamLog({ + credentials, + auditLog + }); + } catch (error) { + logger.error( + error, + `Failed to stream audit log [auditLogId=${auditLog.id}] [provider=${provider}] [orgId=${orgId}]${error instanceof AxiosError ? `: ${error.message}` : ""}` + ); + throw error; + } + }) + ); }; return { @@ -244,6 +217,8 @@ export const auditLogStreamServiceFactory = ({ updateById, deleteById, getById, - list + list, + listProviderOptions, + streamLog }; }; diff --git a/backend/src/ee/services/audit-log-stream/audit-log-stream-types.ts b/backend/src/ee/services/audit-log-stream/audit-log-stream-types.ts index 4c4a5609e..d0f78f632 100644 --- a/backend/src/ee/services/audit-log-stream/audit-log-stream-types.ts +++ b/backend/src/ee/services/audit-log-stream/audit-log-stream-types.ts @@ -1,48 +1,38 @@ -import { TAuditLogStreams } from "@app/db/schemas"; -import { TOrgPermission } from "@app/lib/types"; +import { TAuditLogs } from "@app/db/schemas"; -export type LogStreamHeaders = { - key: string; - value: string; +import { LogProvider } from "./audit-log-stream-enums"; +import { TCustomProvider, TCustomProviderCredentials } from "./custom/custom-provider-types"; +import { TDatadogProvider, TDatadogProviderCredentials } from "./datadog/datadog-provider-types"; +import { TSplunkProvider, TSplunkProviderCredentials } from "./splunk/splunk-provider-types"; + +export type TAuditLogStream = TDatadogProvider | TSplunkProvider | TCustomProvider; + +export type TAuditLogStreamCredentials = + | TDatadogProviderCredentials + | TSplunkProviderCredentials + | TCustomProviderCredentials; + +export type TCreateAuditLogStreamDTO = { + provider: LogProvider; + credentials: TAuditLogStreamCredentials; }; -export type TCreateAuditLogStreamDTO = Omit & { - url: string; - headers?: LogStreamHeaders[]; +export type TUpdateAuditLogStreamDTO = { + logStreamId: string; + provider: LogProvider; + credentials: TAuditLogStreamCredentials; }; -export type TUpdateAuditLogStreamDTO = Omit & { - id: string; - url?: string; - headers?: LogStreamHeaders[]; -}; +export type TLogStreamFactoryValidateCredentials = (input: { + credentials: C; +}) => Promise; -export type TDeleteAuditLogStreamDTO = Omit & { - id: string; -}; +export type TLogStreamFactoryStreamLog = (input: { + credentials: C; + auditLog: TAuditLogs; +}) => Promise; -export type TListAuditLogStreamDTO = Omit; - -export type TGetDetailsAuditLogStreamDTO = Omit & { - id: string; -}; - -export type TAuditLogStreamServiceFactory = { - create: (arg: TCreateAuditLogStreamDTO) => Promise; - updateById: (arg: TUpdateAuditLogStreamDTO) => Promise; - deleteById: (arg: TDeleteAuditLogStreamDTO) => Promise; - getById: (arg: TGetDetailsAuditLogStreamDTO) => Promise<{ - headers: LogStreamHeaders[] | undefined; - orgId: string; - url: string; - id: string; - createdAt: Date; - updatedAt: Date; - encryptedHeadersCiphertext?: string | null | undefined; - encryptedHeadersIV?: string | null | undefined; - encryptedHeadersTag?: string | null | undefined; - encryptedHeadersAlgorithm?: string | null | undefined; - encryptedHeadersKeyEncoding?: string | null | undefined; - }>; - list: (arg: TListAuditLogStreamDTO) => Promise; +export type TLogStreamFactory = () => { + validateCredentials: TLogStreamFactoryValidateCredentials; + streamLog: TLogStreamFactoryStreamLog; }; diff --git a/backend/src/ee/services/audit-log-stream/custom/custom-provider-factory.ts b/backend/src/ee/services/audit-log-stream/custom/custom-provider-factory.ts new file mode 100644 index 000000000..6e397638c --- /dev/null +++ b/backend/src/ee/services/audit-log-stream/custom/custom-provider-factory.ts @@ -0,0 +1,67 @@ +import { RawAxiosRequestHeaders } from "axios"; + +import { request } from "@app/lib/config/request"; +import { BadRequestError } from "@app/lib/errors"; +import { blockLocalAndPrivateIpAddresses } from "@app/lib/validator"; + +import { AUDIT_LOG_STREAM_TIMEOUT } from "../../audit-log/audit-log-queue"; +import { TLogStreamFactoryStreamLog, TLogStreamFactoryValidateCredentials } from "../audit-log-stream-types"; +import { TCustomProviderCredentials } from "./custom-provider-types"; + +export const CustomProviderFactory = () => { + const validateCredentials: TLogStreamFactoryValidateCredentials = async ({ + credentials + }) => { + const { url, headers } = credentials; + + await blockLocalAndPrivateIpAddresses(url); + + const streamHeaders: RawAxiosRequestHeaders = { "Content-Type": "application/json" }; + if (headers.length) { + headers.forEach(({ key, value }) => { + streamHeaders[key] = value; + }); + } + + await request + .post( + url, + { ping: "ok" }, + { + headers: streamHeaders, + timeout: AUDIT_LOG_STREAM_TIMEOUT, + signal: AbortSignal.timeout(AUDIT_LOG_STREAM_TIMEOUT) + } + ) + .catch((err) => { + throw new BadRequestError({ message: `Failed to connect with upstream source: ${(err as Error)?.message}` }); + }); + + return credentials; + }; + + const streamLog: TLogStreamFactoryStreamLog = async ({ credentials, auditLog }) => { + const { url, headers } = credentials; + + await blockLocalAndPrivateIpAddresses(url); + + const streamHeaders: RawAxiosRequestHeaders = { "Content-Type": "application/json" }; + + if (headers.length) { + headers.forEach(({ key, value }) => { + streamHeaders[key] = value; + }); + } + + await request.post(url, auditLog, { + headers: streamHeaders, + timeout: AUDIT_LOG_STREAM_TIMEOUT, + signal: AbortSignal.timeout(AUDIT_LOG_STREAM_TIMEOUT) + }); + }; + + return { + validateCredentials, + streamLog + }; +}; diff --git a/backend/src/ee/services/audit-log-stream/custom/custom-provider-fns.ts b/backend/src/ee/services/audit-log-stream/custom/custom-provider-fns.ts new file mode 100644 index 000000000..27b8bbf72 --- /dev/null +++ b/backend/src/ee/services/audit-log-stream/custom/custom-provider-fns.ts @@ -0,0 +1,8 @@ +import { LogProvider } from "../audit-log-stream-enums"; + +export const getCustomProviderListItem = () => { + return { + name: "Custom" as const, + provider: LogProvider.Custom as const + }; +}; diff --git a/backend/src/ee/services/audit-log-stream/custom/custom-provider-schemas.ts b/backend/src/ee/services/audit-log-stream/custom/custom-provider-schemas.ts new file mode 100644 index 000000000..5d375c42f --- /dev/null +++ b/backend/src/ee/services/audit-log-stream/custom/custom-provider-schemas.ts @@ -0,0 +1,32 @@ +import { z } from "zod"; + +import { LogProvider } from "../audit-log-stream-enums"; + +export const CustomProviderCredentialsSchema = z.object({ + url: z.string().url().trim().min(1).max(255), + headers: z + .object({ + key: z.string().min(1), + value: z.string().min(1) + }) + .array() +}); + +export const CustomProviderSchema = z.object({ + provider: z.literal(LogProvider.Custom), + credentials: CustomProviderCredentialsSchema +}); + +export const SanitizedCustomProviderSchema = z.object({ + provider: z.literal(LogProvider.Custom), + credentials: z.object({ + url: CustomProviderCredentialsSchema.shape.url, + // Only return header keys + headers: CustomProviderCredentialsSchema.shape.headers.element.pick({ key: true }).array() + }) +}); + +export const CustomProviderListItemSchema = z.object({ + name: z.literal("Custom"), + provider: z.literal(LogProvider.Custom) +}); diff --git a/backend/src/ee/services/audit-log-stream/custom/custom-provider-types.ts b/backend/src/ee/services/audit-log-stream/custom/custom-provider-types.ts new file mode 100644 index 000000000..9b2de8347 --- /dev/null +++ b/backend/src/ee/services/audit-log-stream/custom/custom-provider-types.ts @@ -0,0 +1,7 @@ +import { z } from "zod"; + +import { CustomProviderCredentialsSchema, CustomProviderSchema } from "./custom-provider-schemas"; + +export type TCustomProvider = z.infer; + +export type TCustomProviderCredentials = z.infer; diff --git a/backend/src/ee/services/audit-log-stream/datadog/datadog-provider-factory.ts b/backend/src/ee/services/audit-log-stream/datadog/datadog-provider-factory.ts new file mode 100644 index 000000000..f4b8738fa --- /dev/null +++ b/backend/src/ee/services/audit-log-stream/datadog/datadog-provider-factory.ts @@ -0,0 +1,60 @@ +import { RawAxiosRequestHeaders } from "axios"; + +import { request } from "@app/lib/config/request"; +import { BadRequestError } from "@app/lib/errors"; +import { blockLocalAndPrivateIpAddresses } from "@app/lib/validator"; + +import { AUDIT_LOG_STREAM_TIMEOUT } from "../../audit-log/audit-log-queue"; +import { TLogStreamFactoryStreamLog, TLogStreamFactoryValidateCredentials } from "../audit-log-stream-types"; +import { TDatadogProviderCredentials } from "./datadog-provider-types"; + +export const DatadogProviderFactory = () => { + const validateCredentials: TLogStreamFactoryValidateCredentials = async ({ + credentials + }) => { + const { url, token } = credentials; + + await blockLocalAndPrivateIpAddresses(url); + + const streamHeaders: RawAxiosRequestHeaders = { "Content-Type": "application/json", "DD-API-KEY": token }; + + await request + .post( + url, + { ping: "ok" }, + { + headers: streamHeaders, + timeout: AUDIT_LOG_STREAM_TIMEOUT, + signal: AbortSignal.timeout(AUDIT_LOG_STREAM_TIMEOUT) + } + ) + .catch((err) => { + throw new BadRequestError({ message: `Failed to connect with Datadog: ${(err as Error)?.message}` }); + }); + + return credentials; + }; + + const streamLog: TLogStreamFactoryStreamLog = async ({ credentials, auditLog }) => { + const { url, token } = credentials; + + await blockLocalAndPrivateIpAddresses(url); + + const streamHeaders: RawAxiosRequestHeaders = { "Content-Type": "application/json", "DD-API-KEY": token }; + + await request.post( + url, + { ...auditLog, ddsource: "infisical", service: "audit-logs" }, + { + headers: streamHeaders, + timeout: AUDIT_LOG_STREAM_TIMEOUT, + signal: AbortSignal.timeout(AUDIT_LOG_STREAM_TIMEOUT) + } + ); + }; + + return { + validateCredentials, + streamLog + }; +}; diff --git a/backend/src/ee/services/audit-log-stream/datadog/datadog-provider-fns.ts b/backend/src/ee/services/audit-log-stream/datadog/datadog-provider-fns.ts new file mode 100644 index 000000000..ec68fdb39 --- /dev/null +++ b/backend/src/ee/services/audit-log-stream/datadog/datadog-provider-fns.ts @@ -0,0 +1,8 @@ +import { LogProvider } from "../audit-log-stream-enums"; + +export const getDatadogProviderListItem = () => { + return { + name: "Datadog" as const, + provider: LogProvider.Datadog as const + }; +}; diff --git a/backend/src/ee/services/audit-log-stream/datadog/datadog-provider-schemas.ts b/backend/src/ee/services/audit-log-stream/datadog/datadog-provider-schemas.ts new file mode 100644 index 000000000..7adae1ae6 --- /dev/null +++ b/backend/src/ee/services/audit-log-stream/datadog/datadog-provider-schemas.ts @@ -0,0 +1,29 @@ +import RE2 from "re2"; +import { z } from "zod"; + +import { LogProvider } from "../audit-log-stream-enums"; + +export const DatadogProviderCredentialsSchema = z.object({ + url: z.string().url().trim().min(1).max(255), + token: z + .string() + .trim() + .refine((val) => new RE2(/^[a-fA-F0-9]{32}$/).test(val), "Invalid Datadog API key format") +}); + +export const DatadogProviderSchema = z.object({ + provider: z.literal(LogProvider.Datadog), + credentials: DatadogProviderCredentialsSchema +}); + +export const SanitizedDatadogProviderSchema = z.object({ + provider: z.literal(LogProvider.Datadog), + credentials: DatadogProviderCredentialsSchema.pick({ + url: true + }) +}); + +export const DatadogProviderListItemSchema = z.object({ + name: z.literal("Datadog"), + provider: z.literal(LogProvider.Datadog) +}); diff --git a/backend/src/ee/services/audit-log-stream/datadog/datadog-provider-types.ts b/backend/src/ee/services/audit-log-stream/datadog/datadog-provider-types.ts new file mode 100644 index 000000000..e7f6a1f28 --- /dev/null +++ b/backend/src/ee/services/audit-log-stream/datadog/datadog-provider-types.ts @@ -0,0 +1,7 @@ +import { z } from "zod"; + +import { DatadogProviderCredentialsSchema, DatadogProviderSchema } from "./datadog-provider-schemas"; + +export type TDatadogProvider = z.infer; + +export type TDatadogProviderCredentials = z.infer; diff --git a/backend/src/ee/services/audit-log-stream/splunk/splunk-provider-factory.ts b/backend/src/ee/services/audit-log-stream/splunk/splunk-provider-factory.ts new file mode 100644 index 000000000..5481c44eb --- /dev/null +++ b/backend/src/ee/services/audit-log-stream/splunk/splunk-provider-factory.ts @@ -0,0 +1,84 @@ +import { RawAxiosRequestHeaders } from "axios"; + +import { getConfig } from "@app/lib/config/env"; +import { request } from "@app/lib/config/request"; +import { BadRequestError } from "@app/lib/errors"; +import { blockLocalAndPrivateIpAddresses } from "@app/lib/validator"; + +import { AUDIT_LOG_STREAM_TIMEOUT } from "../../audit-log/audit-log-queue"; +import { TLogStreamFactoryStreamLog, TLogStreamFactoryValidateCredentials } from "../audit-log-stream-types"; +import { TSplunkProviderCredentials } from "./splunk-provider-types"; + +function createPayload(event: Record) { + const appCfg = getConfig(); + + return { + time: Math.floor(Date.now() / 1000), + host: new URL(appCfg.SITE_URL || "http://infisical").host, + source: "infisical", + sourcetype: "_json", + event + }; +} + +async function createSplunkUrl(hostname: string) { + let parsedHostname: string; + try { + parsedHostname = new URL(`https://${hostname}`).hostname; + } catch (error) { + throw new BadRequestError({ message: `Invalid Splunk hostname provided: ${(error as Error).message}` }); + } + + await blockLocalAndPrivateIpAddresses(`https://${hostname}`); + + return `https://${parsedHostname}:8088/services/collector/event`; +} + +export const SplunkProviderFactory = () => { + const validateCredentials: TLogStreamFactoryValidateCredentials = async ({ + credentials + }) => { + const { hostname, token } = credentials; + + const url = await createSplunkUrl(hostname); + + const streamHeaders: RawAxiosRequestHeaders = { + "Content-Type": "application/json", + Authorization: `Splunk ${token}` + }; + + await request + .post(url, createPayload({ ping: "ok" }), { + headers: streamHeaders, + timeout: AUDIT_LOG_STREAM_TIMEOUT, + signal: AbortSignal.timeout(AUDIT_LOG_STREAM_TIMEOUT) + }) + .catch((err) => { + throw new BadRequestError({ message: `Failed to connect with Splunk: ${(err as Error)?.message}` }); + }); + + return credentials; + }; + + const streamLog: TLogStreamFactoryStreamLog = async ({ credentials, auditLog }) => { + const { hostname, token } = credentials; + + const url = await createSplunkUrl(hostname); + + const streamHeaders: RawAxiosRequestHeaders = { + "Content-Type": "application/json", + Authorization: `Splunk ${token}` + }; + + await request.post(url, createPayload(auditLog), { + headers: streamHeaders, + timeout: AUDIT_LOG_STREAM_TIMEOUT, + signal: AbortSignal.timeout(AUDIT_LOG_STREAM_TIMEOUT) + }); + }; + + return { + validateCredentials, + streamLog + }; +}; diff --git a/backend/src/ee/services/audit-log-stream/splunk/splunk-provider-fns.ts b/backend/src/ee/services/audit-log-stream/splunk/splunk-provider-fns.ts new file mode 100644 index 000000000..e2ea2e316 --- /dev/null +++ b/backend/src/ee/services/audit-log-stream/splunk/splunk-provider-fns.ts @@ -0,0 +1,8 @@ +import { LogProvider } from "../audit-log-stream-enums"; + +export const getSplunkProviderListItem = () => { + return { + name: "Splunk" as const, + provider: LogProvider.Splunk as const + }; +}; diff --git a/backend/src/ee/services/audit-log-stream/splunk/splunk-provider-schemas.ts b/backend/src/ee/services/audit-log-stream/splunk/splunk-provider-schemas.ts new file mode 100644 index 000000000..b03524951 --- /dev/null +++ b/backend/src/ee/services/audit-log-stream/splunk/splunk-provider-schemas.ts @@ -0,0 +1,25 @@ +import { z } from "zod"; + +import { LogProvider } from "../audit-log-stream-enums"; + +export const SplunkProviderCredentialsSchema = z.object({ + hostname: z.string().url().trim().min(1).max(255), + token: z.string().uuid().trim().min(1) +}); + +export const SplunkProviderSchema = z.object({ + provider: z.literal(LogProvider.Splunk), + credentials: SplunkProviderCredentialsSchema +}); + +export const SanitizedSplunkProviderSchema = z.object({ + provider: z.literal(LogProvider.Splunk), + credentials: SplunkProviderCredentialsSchema.pick({ + hostname: true + }) +}); + +export const SplunkProviderListItemSchema = z.object({ + name: z.literal("Splunk"), + provider: z.literal(LogProvider.Splunk) +}); diff --git a/backend/src/ee/services/audit-log-stream/splunk/splunk-provider-types.ts b/backend/src/ee/services/audit-log-stream/splunk/splunk-provider-types.ts new file mode 100644 index 000000000..11f7c8fcb --- /dev/null +++ b/backend/src/ee/services/audit-log-stream/splunk/splunk-provider-types.ts @@ -0,0 +1,7 @@ +import { z } from "zod"; + +import { SplunkProviderCredentialsSchema, SplunkProviderSchema } from "./splunk-provider-schemas"; + +export type TSplunkProvider = z.infer; + +export type TSplunkProviderCredentials = z.infer; diff --git a/backend/src/ee/services/audit-log/audit-log-queue.ts b/backend/src/ee/services/audit-log/audit-log-queue.ts index 0914b8f6b..6b286c2ec 100644 --- a/backend/src/ee/services/audit-log/audit-log-queue.ts +++ b/backend/src/ee/services/audit-log/audit-log-queue.ts @@ -1,22 +1,14 @@ -import { AxiosError, RawAxiosRequestHeaders } from "axios"; - -import { SecretKeyEncoding } from "@app/db/schemas"; -import { request } from "@app/lib/config/request"; -import { crypto } from "@app/lib/crypto/cryptography"; -import { logger } from "@app/lib/logger"; +import { TAuditLogStreamServiceFactory } from "@app/ee/services/audit-log-stream/audit-log-stream-service"; import { QueueJobs, QueueName, TQueueServiceFactory } from "@app/queue"; import { TProjectDALFactory } from "@app/services/project/project-dal"; -import { TAuditLogStreamDALFactory } from "../audit-log-stream/audit-log-stream-dal"; -import { providerSpecificPayload } from "../audit-log-stream/audit-log-stream-fns"; -import { LogStreamHeaders } from "../audit-log-stream/audit-log-stream-types"; import { TLicenseServiceFactory } from "../license/license-service"; import { TAuditLogDALFactory } from "./audit-log-dal"; import { TCreateAuditLogDTO } from "./audit-log-types"; type TAuditLogQueueServiceFactoryDep = { auditLogDAL: TAuditLogDALFactory; - auditLogStreamDAL: Pick; + auditLogStreamService: Pick; queueService: TQueueServiceFactory; projectDAL: Pick; licenseService: Pick; @@ -35,7 +27,7 @@ export const auditLogQueueServiceFactory = async ({ queueService, projectDAL, licenseService, - auditLogStreamDAL + auditLogStreamService }: TAuditLogQueueServiceFactoryDep): Promise => { const pushToLog = async (data: TCreateAuditLogDTO) => { await queueService.queue(QueueName.AuditLog, QueueJobs.AuditLog, data, { @@ -86,60 +78,7 @@ export const auditLogQueueServiceFactory = async ({ userAgentType }); - const logStreams = orgId ? await auditLogStreamDAL.find({ orgId }) : []; - await Promise.allSettled( - logStreams.map( - async ({ - url, - encryptedHeadersTag, - encryptedHeadersIV, - encryptedHeadersKeyEncoding, - encryptedHeadersCiphertext - }) => { - const streamHeaders = - encryptedHeadersIV && encryptedHeadersCiphertext && encryptedHeadersTag - ? (JSON.parse( - crypto - .encryption() - .symmetric() - .decryptWithRootEncryptionKey({ - keyEncoding: encryptedHeadersKeyEncoding as SecretKeyEncoding, - iv: encryptedHeadersIV, - tag: encryptedHeadersTag, - ciphertext: encryptedHeadersCiphertext - }) - ) as LogStreamHeaders[]) - : []; - - const headers: RawAxiosRequestHeaders = { "Content-Type": "application/json" }; - - if (streamHeaders.length) - streamHeaders.forEach(({ key, value }) => { - headers[key] = value; - }); - - try { - const response = await request.post( - url, - { ...providerSpecificPayload(url), ...auditLog }, - { - headers, - // request timeout - timeout: AUDIT_LOG_STREAM_TIMEOUT, - // connection timeout - signal: AbortSignal.timeout(AUDIT_LOG_STREAM_TIMEOUT) - } - ); - return response; - } catch (error) { - logger.error( - `Failed to stream audit log [url=${url}] for org [orgId=${orgId}] [error=${(error as AxiosError).message}]` - ); - return error; - } - } - ) - ); + await auditLogStreamService.streamLog(orgId, auditLog); } }); diff --git a/backend/src/server/routes/index.ts b/backend/src/server/routes/index.ts index 3773048b0..7c6f47b60 100644 --- a/backend/src/server/routes/index.ts +++ b/backend/src/server/routes/index.ts @@ -555,20 +555,22 @@ export const registerRoutes = async ( permissionService }); + const auditLogStreamService = auditLogStreamServiceFactory({ + licenseService, + permissionService, + auditLogStreamDAL, + kmsService + }); + const auditLogQueue = await auditLogQueueServiceFactory({ auditLogDAL, queueService, projectDAL, licenseService, - auditLogStreamDAL + auditLogStreamService }); const auditLogService = auditLogServiceFactory({ auditLogDAL, permissionService, auditLogQueue }); - const auditLogStreamService = auditLogStreamServiceFactory({ - licenseService, - permissionService, - auditLogStreamDAL - }); const secretApprovalPolicyService = secretApprovalPolicyServiceFactory({ projectEnvDAL, secretApprovalPolicyApproverDAL: sapApproverDAL, diff --git a/backend/src/server/routes/sanitizedSchemas.ts b/backend/src/server/routes/sanitizedSchemas.ts index 8af4baa8b..344224008 100644 --- a/backend/src/server/routes/sanitizedSchemas.ts +++ b/backend/src/server/routes/sanitizedSchemas.ts @@ -246,13 +246,6 @@ export const SanitizedDynamicSecretSchema = DynamicSecretsSchema.omit({ metadata: ResourceMetadataSchema.optional() }); -export const SanitizedAuditLogStreamSchema = z.object({ - id: z.string(), - url: z.string(), - createdAt: z.date(), - updatedAt: z.date() -}); - export const SanitizedProjectSchema = ProjectsSchema.pick({ id: true, name: true,