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 new file mode 100644 index 000000000..dc93f238e --- /dev/null +++ b/backend/src/ee/services/audit-log-stream/audit-log-stream-fns.ts @@ -0,0 +1,21 @@ +export function providerSpecificPayload(url: string) { + const { hostname } = new URL(url); + + const payload: Record = {}; + + 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; + } + + return payload; +} 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 72207aca1..65e49bdea 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 @@ -13,6 +13,7 @@ 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"; type TAuditLogStreamServiceFactoryDep = { @@ -69,10 +70,11 @@ export const auditLogStreamServiceFactory = ({ headers.forEach(({ key, value }) => { streamHeaders[key] = value; }); + await request .post( url, - { ping: "ok" }, + { ...providerSpecificPayload(url), ping: "ok" }, { headers: streamHeaders, // request timeout @@ -137,7 +139,7 @@ export const auditLogStreamServiceFactory = ({ await request .post( url || logStream.url, - { ping: "ok" }, + { ...providerSpecificPayload(url || logStream.url), ping: "ok" }, { headers: streamHeaders, // request timeout 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 0f774e911..e812cb7ea 100644 --- a/backend/src/ee/services/audit-log/audit-log-queue.ts +++ b/backend/src/ee/services/audit-log/audit-log-queue.ts @@ -1,13 +1,15 @@ -import { RawAxiosRequestHeaders } from "axios"; +import { AxiosError, RawAxiosRequestHeaders } from "axios"; import { SecretKeyEncoding } from "@app/db/schemas"; import { getConfig } from "@app/lib/config/env"; import { request } from "@app/lib/config/request"; import { infisicalSymmetricDecrypt } from "@app/lib/crypto/encryption"; +import { logger } from "@app/lib/logger"; 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"; @@ -128,13 +130,29 @@ export const auditLogQueueServiceFactory = async ({ headers[key] = value; }); - return request.post(url, auditLog, { - headers, - // request timeout - timeout: AUDIT_LOG_STREAM_TIMEOUT, - // connection timeout - signal: AbortSignal.timeout(AUDIT_LOG_STREAM_TIMEOUT) - }); + try { + logger.info(`Streaming audit log [url=${url}] for org [orgId=${orgId}]`); + 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) + } + ); + logger.info( + `Successfully streamed audit log [url=${url}] for org [orgId=${orgId}] [response=${JSON.stringify(response.data)}]` + ); + return response; + } catch (error) { + logger.error( + `Failed to stream audit log [url=${url}] for org [orgId=${orgId}] [error=${(error as AxiosError).message}]` + ); + return error; + } } ) ); @@ -218,13 +236,29 @@ export const auditLogQueueServiceFactory = async ({ headers[key] = value; }); - return request.post(url, auditLog, { - headers, - // request timeout - timeout: AUDIT_LOG_STREAM_TIMEOUT, - // connection timeout - signal: AbortSignal.timeout(AUDIT_LOG_STREAM_TIMEOUT) - }); + try { + logger.info(`Streaming audit log [url=${url}] for org [orgId=${orgId}]`); + 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) + } + ); + logger.info( + `Successfully streamed audit log [url=${url}] for org [orgId=${orgId}] [response=${JSON.stringify(response.data)}]` + ); + return response; + } catch (error) { + logger.error( + `Failed to stream audit log [url=${url}] for org [orgId=${orgId}] [error=${(error as AxiosError).message}]` + ); + return error; + } } ) );