diff --git a/openframe-data-cassandra/src/main/java/com/openframe/data/cassandra/model/enums/UnifiedEventType.java b/openframe-data-cassandra/src/main/java/com/openframe/data/cassandra/model/enums/UnifiedEventType.java index f5927928e..3e87c12a0 100644 --- a/openframe-data-cassandra/src/main/java/com/openframe/data/cassandra/model/enums/UnifiedEventType.java +++ b/openframe-data-cassandra/src/main/java/com/openframe/data/cassandra/model/enums/UnifiedEventType.java @@ -164,6 +164,25 @@ public enum UnifiedEventType { INTEGRATION_UPDATED(Severity.INFO, "Integration updated"), INTEGRATION_DELETED(Severity.INFO, "Integration deleted"), + // Microsoft 365 directory audit events + M365_USER_MANAGEMENT(Severity.INFO, "Microsoft 365 user management"), + M365_GROUP_MANAGEMENT(Severity.INFO, "Microsoft 365 group management"), + M365_APPLICATION_MANAGEMENT(Severity.INFO, "Microsoft 365 application management"), + M365_ROLE_MANAGEMENT(Severity.WARNING, "Microsoft 365 role management"), + M365_POLICY(Severity.INFO, "Microsoft 365 policy change"), + M365_DIRECTORY_MANAGEMENT(Severity.INFO, "Microsoft 365 directory management"), + M365_AUDIT_FAILURE(Severity.WARNING, "Microsoft 365 failed directory operation"), + M365_AUDIT_OTHER(Severity.INFO, "Microsoft 365 directory audit event"), + + // Google Workspace directory audit events + GWS_USER_MANAGEMENT(Severity.INFO, "Google Workspace user management"), + GWS_GROUP_MANAGEMENT(Severity.INFO, "Google Workspace group management"), + GWS_SECURITY_SETTINGS(Severity.WARNING, "Google Workspace security settings change"), + GWS_DOMAIN_SETTINGS(Severity.INFO, "Google Workspace domain settings change"), + GWS_DELEGATED_ADMIN(Severity.WARNING, "Google Workspace delegated admin change"), + GWS_AUDIT_FAILURE(Severity.WARNING, "Google Workspace failed directory operation"), + GWS_AUDIT_OTHER(Severity.INFO, "Google Workspace directory audit event"), + // Unknown events UNKNOWN(Severity.WARNING, "Unknown event"); diff --git a/openframe-data-kafka/src/main/java/com/openframe/data/model/enums/DataEnrichmentServiceType.java b/openframe-data-kafka/src/main/java/com/openframe/data/model/enums/DataEnrichmentServiceType.java index 02e94179b..087ecc220 100644 --- a/openframe-data-kafka/src/main/java/com/openframe/data/model/enums/DataEnrichmentServiceType.java +++ b/openframe-data-kafka/src/main/java/com/openframe/data/model/enums/DataEnrichmentServiceType.java @@ -3,6 +3,7 @@ public enum DataEnrichmentServiceType { INTEGRATED_TOOLS_EVENTS, - RMM_RESULTS + RMM_RESULTS, + PRE_ENRICHED } diff --git a/openframe-data-kafka/src/main/java/com/openframe/data/model/enums/IntegratedToolType.java b/openframe-data-kafka/src/main/java/com/openframe/data/model/enums/IntegratedToolType.java index 498e4cd5c..8bc4b79b7 100644 --- a/openframe-data-kafka/src/main/java/com/openframe/data/model/enums/IntegratedToolType.java +++ b/openframe-data-kafka/src/main/java/com/openframe/data/model/enums/IntegratedToolType.java @@ -4,7 +4,9 @@ public enum IntegratedToolType { RMM("rmm"), MESHCENTRAL ("meshcentral"), - FLEET ("fleet-mdm"); + FLEET ("fleet-mdm"), + MICROSOFT_365("microsoft-365"), + GOOGLE_WORKSPACE("google-workspace"); private final String dbName; diff --git a/openframe-data-kafka/src/main/java/com/openframe/data/model/enums/MessageType.java b/openframe-data-kafka/src/main/java/com/openframe/data/model/enums/MessageType.java index 02d3947f6..fa59b6ad7 100644 --- a/openframe-data-kafka/src/main/java/com/openframe/data/model/enums/MessageType.java +++ b/openframe-data-kafka/src/main/java/com/openframe/data/model/enums/MessageType.java @@ -20,6 +20,10 @@ public enum MessageType { FLEET_MDM_POLICY_ACTIVITY_EVENT(IntegratedToolType.FLEET, DataEnrichmentServiceType.INTEGRATED_TOOLS_EVENTS, List.of(Destination.CASSANDRA_EVENT_LOG, Destination.KAFKA_PINOT), EventHandlerType.COMMON_TYPE), FLEET_MDM_POLICY_MEMBERSHIP_EVENT(IntegratedToolType.FLEET, DataEnrichmentServiceType.INTEGRATED_TOOLS_EVENTS, + List.of(Destination.CASSANDRA_EVENT_LOG, Destination.KAFKA_PINOT), EventHandlerType.COMMON_TYPE), + MICROSOFT_365_AUDIT_EVENT(IntegratedToolType.MICROSOFT_365, DataEnrichmentServiceType.PRE_ENRICHED, + List.of(Destination.CASSANDRA_EVENT_LOG, Destination.KAFKA_PINOT), EventHandlerType.COMMON_TYPE), + GOOGLE_WORKSPACE_AUDIT_EVENT(IntegratedToolType.GOOGLE_WORKSPACE, DataEnrichmentServiceType.PRE_ENRICHED, List.of(Destination.CASSANDRA_EVENT_LOG, Destination.KAFKA_PINOT), EventHandlerType.COMMON_TYPE); private final IntegratedToolType integratedToolType; diff --git a/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/GoogleWorkspaceAuditEventDeserializer.java b/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/GoogleWorkspaceAuditEventDeserializer.java new file mode 100644 index 000000000..a46c3565a --- /dev/null +++ b/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/GoogleWorkspaceAuditEventDeserializer.java @@ -0,0 +1,145 @@ +package com.openframe.stream.deserializer; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.fasterxml.jackson.databind.node.ObjectNode; +import com.openframe.data.cassandra.model.enums.UnifiedEventType; +import com.openframe.data.model.enums.MessageType; +import com.openframe.kafka.model.debezium.CommonDebeziumMessage; +import com.openframe.stream.mapping.EventTypeMapper; +import com.openframe.stream.model.fleet.debezium.DeserializedDebeziumMessage; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang3.StringUtils; +import org.springframework.stereotype.Component; + +import java.time.Instant; +import java.time.ZoneId; +import java.time.format.DateTimeFormatter; +import java.util.Optional; + +/** + * Deserializes Google Workspace directory audit events polled from the Admin SDK Reports API + * {@code activities.list("admin")} endpoint. Events are hand-built by the poller (not CDC): a + * single polled activity's {@code events[]} array is fanned out by the poller into one Kafka + * record per event, so this deserializer stays 1:1 with {@link KafkaMessageDeserializer#deserialize} + * (that interface is single-result). The poller writes a FLAT {@code after} object per event with + * fields {@code uniqueQualifier}, {@code eventIndex}, {@code activityTime}, {@code eventType}, + * {@code eventName}, {@code actorEmail}, {@code ipAddress}, {@code event} (nested JSON carrying the + * raw event's {@code parameters}), plus tenant/organization passthrough fields ({@code tenantId}, + * {@code organizationId}, {@code organizationName}) and multi-connection fields + * ({@code connectionId}, {@code connectionName}) — hence + * {@link com.openframe.data.model.enums.DataEnrichmentServiceType#PRE_ENRICHED}. {@code toolEventId} + * is {@code uniqueQualifier + "-" + eventIndex}: Reports API activities are uniquely identified by + * {@code uniqueQualifier}, but a single activity can carry multiple events, so the pair keeps + * replays from the poller's cursor overlap window upsert idempotent per event. Events carry no + * agent reference. + * {@code connectionId}/{@code connectionName} (multi-connection orgs) are passed through into details. + *

+ * {@code event.parameters[]} is an UNTOUCHED passthrough of the Reports API parameter union: each + * entry always has {@code name}, but its value key varies — {@code value} (string), + * {@code boolValue}, {@code intValue}, {@code multiValue} (array of strings) or + * {@code multiMessageValue}. The write-audit publisher (saas-lib + * {@code GoogleWorkspaceWriteAuditPublisher}) emits only the string {@code {name,value}} member of + * that union. Consumers rendering details must read + * {@code value ?? boolValue ?? intValue ?? multiValue} for the scalar/array members; nothing may + * assume {@code value} alone. {@code multiMessageValue} is deliberately NOT part of that fallback + * chain — it is an array of nested {@code {parameter: [{name, value, ...}]}} objects, so a consumer + * that needs it must render it recursively rather than coerce it to a scalar. + */ +@Slf4j +@Component +@RequiredArgsConstructor +public class GoogleWorkspaceAuditEventDeserializer implements KafkaMessageDeserializer { + + private static final DateTimeFormatter DAY_FORMATTER = + DateTimeFormatter.ofPattern("yyyy-MM-dd").withZone(ZoneId.of("UTC")); + // Reports API does not carry a structured success/failure result field; failure is signalled + // by eventName itself (e.g. LOGIN_FAILURE, login_failure). + private static final String FAILURE_MARKER = "_FAILURE"; + private static final String UNKNOWN = "unknown"; + + private final ObjectMapper mapper; + + @Override + public MessageType getType() { + return MessageType.GOOGLE_WORKSPACE_AUDIT_EVENT; + } + + @Override + public DeserializedDebeziumMessage deserialize(CommonDebeziumMessage debeziumMessage, MessageType messageType) { + JsonNode after = debeziumMessage.getPayload().getAfter(); + if (after == null || after.isNull()) { + return null; + } + long eventTimestamp = getEventTimestamp(after) + .orElse(debeziumMessage.getPayload().getTimestamp()); + String eventType = textField(after, "eventType").orElse(UNKNOWN); + + return DeserializedDebeziumMessage.builder() + .payload(debeziumMessage.getPayload()) + .agentId(null) + .ingestDay(DAY_FORMATTER.format(Instant.ofEpochMilli(eventTimestamp))) + .sourceEventType(eventType) + .toolEventId(buildToolEventId(after)) + .unifiedEventType(resolveEventType(after, eventType)) + .message(textField(after, "eventName").orElse(null)) + .integratedToolType(messageType.getIntegratedToolType()) + .debeziumMessage(after.toString()) + .details(buildDetails(after)) + .eventTimestamp(eventTimestamp) + .skipProcessing(false) + .isVisible(true) + .tenantId(textField(after, "tenantId").orElse(null)) + .organizationId(textField(after, "organizationId").orElse(null)) + .organizationName(textField(after, "organizationName").orElse(null)) + .userId(textField(after, "actorEmail").orElse(null)) + .build(); + } + + private UnifiedEventType resolveEventType(JsonNode after, String eventType) { + String eventName = textField(after, "eventName").orElse(""); + if (StringUtils.containsIgnoreCase(eventName, FAILURE_MARKER)) { + return UnifiedEventType.GWS_AUDIT_FAILURE; + } + UnifiedEventType mapped = EventTypeMapper.mapToUnifiedType(getType().getIntegratedToolType(), eventType); + return mapped == UnifiedEventType.UNKNOWN ? UnifiedEventType.GWS_AUDIT_OTHER : mapped; + } + + private String buildToolEventId(JsonNode after) { + return textField(after, "uniqueQualifier") + .map(uniqueQualifier -> uniqueQualifier + "-" + textField(after, "eventIndex").orElse("0")) + .orElse(null); + } + + private Optional getEventTimestamp(JsonNode after) { + return textField(after, "activityTime") + .flatMap(value -> { + try { + return Optional.of(Instant.parse(value).toEpochMilli()); + } catch (Exception e) { + log.warn("Unparseable activityTime '{}', falling back to processing timestamp", value); + return Optional.empty(); + } + }); + } + + private String buildDetails(JsonNode after) { + ObjectNode details = mapper.createObjectNode(); + JsonNode event = after.get("event"); + if (event != null && !event.isNull()) { + details.set("event", event); + } + textField(after, "ipAddress").ifPresent(value -> details.put("ipAddress", value)); + textField(after, "connectionId").ifPresent(value -> details.put("connectionId", value)); + textField(after, "connectionName").ifPresent(value -> details.put("connectionName", value)); + return details.toString(); + } + + private Optional textField(JsonNode node, String fieldName) { + return Optional.ofNullable(node.get(fieldName)) + .filter(field -> !field.isNull()) + .map(JsonNode::asText) + .filter(StringUtils::isNotBlank); + } +} diff --git a/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/Microsoft365AuditEventDeserializer.java b/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/Microsoft365AuditEventDeserializer.java new file mode 100644 index 000000000..2e478de1d --- /dev/null +++ b/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/Microsoft365AuditEventDeserializer.java @@ -0,0 +1,131 @@ +package com.openframe.stream.deserializer; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.fasterxml.jackson.databind.node.ObjectNode; +import com.openframe.data.cassandra.model.enums.UnifiedEventType; +import com.openframe.data.model.enums.MessageType; +import com.openframe.kafka.model.debezium.CommonDebeziumMessage; +import com.openframe.stream.mapping.EventTypeMapper; +import com.openframe.stream.model.fleet.debezium.DeserializedDebeziumMessage; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang3.StringUtils; +import org.springframework.stereotype.Component; + +import java.time.Instant; +import java.time.ZoneId; +import java.time.format.DateTimeFormatter; +import java.util.Optional; +import java.util.Set; + +/** + * Deserializes Microsoft 365 Entra directory audit events polled from Graph + * {@code auditLogs/directoryAudits}. Events are hand-built by the poller (not CDC), arrive + * pre-enriched with tenant/organization fields in the payload and carry no agent reference — + * hence {@link com.openframe.data.model.enums.DataEnrichmentServiceType#PRE_ENRICHED}. + * {@code toolEventId} is the Graph audit record id, so replays from the poller's cursor + * overlap window upsert idempotently. + * {@code connectionId}/{@code connectionName} (multi-connection orgs) are passed through into details. + */ +@Slf4j +@Component +@RequiredArgsConstructor +public class Microsoft365AuditEventDeserializer implements KafkaMessageDeserializer { + + private static final DateTimeFormatter DAY_FORMATTER = + DateTimeFormatter.ofPattern("yyyy-MM-dd").withZone(ZoneId.of("UTC")); + // Graph result values that mean the operation did not succeed (result can be + // success | failure | timeout | unknownFutureValue). + private static final Set FAILED_RESULTS = Set.of("failure", "timeout"); + private static final String UNKNOWN = "unknown"; + + private final ObjectMapper mapper; + + @Override + public MessageType getType() { + return MessageType.MICROSOFT_365_AUDIT_EVENT; + } + + @Override + public DeserializedDebeziumMessage deserialize(CommonDebeziumMessage debeziumMessage, MessageType messageType) { + JsonNode after = debeziumMessage.getPayload().getAfter(); + if (after == null || after.isNull()) { + return null; + } + long eventTimestamp = getEventTimestamp(after) + .orElse(debeziumMessage.getPayload().getTimestamp()); + String category = textField(after, "category").orElse(UNKNOWN); + + return DeserializedDebeziumMessage.builder() + .payload(debeziumMessage.getPayload()) + .agentId(null) + .ingestDay(DAY_FORMATTER.format(Instant.ofEpochMilli(eventTimestamp))) + .sourceEventType(category) + .toolEventId(textField(after, "auditId").orElse(null)) + .unifiedEventType(resolveEventType(after, category)) + .message(textField(after, "activityDisplayName").orElse(null)) + .integratedToolType(messageType.getIntegratedToolType()) + .debeziumMessage(after.toString()) + .details(buildDetails(after)) + .eventTimestamp(eventTimestamp) + .skipProcessing(false) + .isVisible(true) + .tenantId(textField(after, "tenantId").orElse(null)) + .organizationId(textField(after, "organizationId").orElse(null)) + .organizationName(textField(after, "organizationName").orElse(null)) + .userId(textPath(after.path("initiatedBy").path("user"), "userPrincipalName")) + .build(); + } + + private UnifiedEventType resolveEventType(JsonNode after, String category) { + if (textField(after, "result").map(String::toLowerCase).filter(FAILED_RESULTS::contains).isPresent()) { + return UnifiedEventType.M365_AUDIT_FAILURE; + } + UnifiedEventType mapped = EventTypeMapper.mapToUnifiedType(getType().getIntegratedToolType(), category); + return mapped == UnifiedEventType.UNKNOWN ? UnifiedEventType.M365_AUDIT_OTHER : mapped; + } + + private Optional getEventTimestamp(JsonNode after) { + return textField(after, "activityDateTime") + .flatMap(value -> { + try { + return Optional.of(Instant.parse(value).toEpochMilli()); + } catch (Exception e) { + log.warn("Unparseable activityDateTime '{}', falling back to processing timestamp", value); + return Optional.empty(); + } + }); + } + + private String buildDetails(JsonNode after) { + ObjectNode details = mapper.createObjectNode(); + JsonNode initiatedBy = after.get("initiatedBy"); + if (initiatedBy != null && !initiatedBy.isNull()) { + details.set("initiatedBy", initiatedBy); + } + JsonNode targetResources = after.get("targetResources"); + if (targetResources != null && !targetResources.isNull()) { + details.set("targetResources", targetResources); + } + JsonNode additionalDetails = after.get("additionalDetails"); + if (additionalDetails != null && !additionalDetails.isNull()) { + details.set("additionalDetails", additionalDetails); + } + textField(after, "connectionId").ifPresent(value -> details.put("connectionId", value)); + textField(after, "connectionName").ifPresent(value -> details.put("connectionName", value)); + return details.toString(); + } + + private Optional textField(JsonNode node, String fieldName) { + return Optional.ofNullable(node.get(fieldName)) + .filter(field -> !field.isNull()) + .map(JsonNode::asText) + .filter(StringUtils::isNotBlank); + } + + private String textPath(JsonNode node, String fieldName) { + String value = node.path(fieldName).asText(null); + return StringUtils.isNotBlank(value) ? value : null; + } +} diff --git a/openframe-stream-service-core/src/main/java/com/openframe/stream/mapping/EventTypeMapper.java b/openframe-stream-service-core/src/main/java/com/openframe/stream/mapping/EventTypeMapper.java index bc58d0d98..5474661e5 100644 --- a/openframe-stream-service-core/src/main/java/com/openframe/stream/mapping/EventTypeMapper.java +++ b/openframe-stream-service-core/src/main/java/com/openframe/stream/mapping/EventTypeMapper.java @@ -269,5 +269,20 @@ private static void initializeDefaultMappings() { // Policy Membership registerMapping(IntegratedToolType.FLEET, SourceEventTypes.Fleet.POLICY_MEMBERSHIP_PASS, UnifiedEventType.POLICY_APPLIED); registerMapping(IntegratedToolType.FLEET, SourceEventTypes.Fleet.POLICY_MEMBERSHIP_FAIL, UnifiedEventType.POLICY_VIOLATION); + + // Microsoft 365 Entra directory audit mappings (directoryAudits category values) + registerMapping(IntegratedToolType.MICROSOFT_365, SourceEventTypes.Microsoft365.USER_MANAGEMENT, UnifiedEventType.M365_USER_MANAGEMENT); + registerMapping(IntegratedToolType.MICROSOFT_365, SourceEventTypes.Microsoft365.GROUP_MANAGEMENT, UnifiedEventType.M365_GROUP_MANAGEMENT); + registerMapping(IntegratedToolType.MICROSOFT_365, SourceEventTypes.Microsoft365.APPLICATION_MANAGEMENT, UnifiedEventType.M365_APPLICATION_MANAGEMENT); + registerMapping(IntegratedToolType.MICROSOFT_365, SourceEventTypes.Microsoft365.ROLE_MANAGEMENT, UnifiedEventType.M365_ROLE_MANAGEMENT); + registerMapping(IntegratedToolType.MICROSOFT_365, SourceEventTypes.Microsoft365.POLICY, UnifiedEventType.M365_POLICY); + registerMapping(IntegratedToolType.MICROSOFT_365, SourceEventTypes.Microsoft365.DIRECTORY_MANAGEMENT, UnifiedEventType.M365_DIRECTORY_MANAGEMENT); + + // Google Workspace directory audit mappings (admin application event types) + registerMapping(IntegratedToolType.GOOGLE_WORKSPACE, SourceEventTypes.GoogleWorkspace.USER_SETTINGS, UnifiedEventType.GWS_USER_MANAGEMENT); + registerMapping(IntegratedToolType.GOOGLE_WORKSPACE, SourceEventTypes.GoogleWorkspace.GROUP_SETTINGS, UnifiedEventType.GWS_GROUP_MANAGEMENT); + registerMapping(IntegratedToolType.GOOGLE_WORKSPACE, SourceEventTypes.GoogleWorkspace.SECURITY_SETTINGS, UnifiedEventType.GWS_SECURITY_SETTINGS); + registerMapping(IntegratedToolType.GOOGLE_WORKSPACE, SourceEventTypes.GoogleWorkspace.DOMAIN_SETTINGS, UnifiedEventType.GWS_DOMAIN_SETTINGS); + registerMapping(IntegratedToolType.GOOGLE_WORKSPACE, SourceEventTypes.GoogleWorkspace.DELEGATED_ADMIN_SETTINGS, UnifiedEventType.GWS_DELEGATED_ADMIN); } } diff --git a/openframe-stream-service-core/src/main/java/com/openframe/stream/mapping/SourceEventTypes.java b/openframe-stream-service-core/src/main/java/com/openframe/stream/mapping/SourceEventTypes.java index a826c4936..f3655b46b 100644 --- a/openframe-stream-service-core/src/main/java/com/openframe/stream/mapping/SourceEventTypes.java +++ b/openframe-stream-service-core/src/main/java/com/openframe/stream/mapping/SourceEventTypes.java @@ -250,4 +250,29 @@ interface Fleet { String POLICY_MEMBERSHIP_PASS = "policy_membership_pass"; String POLICY_MEMBERSHIP_FAIL = "policy_membership_fail"; } + + /** + * Microsoft 365 Entra directory audit event types (Graph directoryAudits {@code category} values). + */ + interface Microsoft365 { + + String USER_MANAGEMENT = "UserManagement"; + String GROUP_MANAGEMENT = "GroupManagement"; + String APPLICATION_MANAGEMENT = "ApplicationManagement"; + String ROLE_MANAGEMENT = "RoleManagement"; + String POLICY = "Policy"; + String DIRECTORY_MANAGEMENT = "DirectoryManagement"; + } + + /** + * Google Workspace directory audit event types (Reports API {@code admin} application event types). + */ + interface GoogleWorkspace { + + String USER_SETTINGS = "USER_SETTINGS"; + String GROUP_SETTINGS = "GROUP_SETTINGS"; + String SECURITY_SETTINGS = "SECURITY_SETTINGS"; + String DOMAIN_SETTINGS = "DOMAIN_SETTINGS"; + String DELEGATED_ADMIN_SETTINGS = "DELEGATED_ADMIN_SETTINGS"; + } } diff --git a/openframe-stream-service-core/src/main/java/com/openframe/stream/model/fleet/debezium/DeserializedDebeziumMessage.java b/openframe-stream-service-core/src/main/java/com/openframe/stream/model/fleet/debezium/DeserializedDebeziumMessage.java index 1a80140f7..462ff8b7d 100644 --- a/openframe-stream-service-core/src/main/java/com/openframe/stream/model/fleet/debezium/DeserializedDebeziumMessage.java +++ b/openframe-stream-service-core/src/main/java/com/openframe/stream/model/fleet/debezium/DeserializedDebeziumMessage.java @@ -28,5 +28,8 @@ public class DeserializedDebeziumMessage extends CommonDebeziumMessage { private Boolean skipProcessing; private Boolean isVisible; private String tenantId; + private String organizationId; + private String organizationName; + private String userId; private Set excludedDestinations; } diff --git a/openframe-stream-service-core/src/main/java/com/openframe/stream/service/PreEnrichedDataEnrichmentService.java b/openframe-stream-service-core/src/main/java/com/openframe/stream/service/PreEnrichedDataEnrichmentService.java new file mode 100644 index 000000000..7b9ce1665 --- /dev/null +++ b/openframe-stream-service-core/src/main/java/com/openframe/stream/service/PreEnrichedDataEnrichmentService.java @@ -0,0 +1,47 @@ +package com.openframe.stream.service; + +import com.openframe.data.model.enums.DataEnrichmentServiceType; +import com.openframe.data.service.TenantIdProvider; +import com.openframe.stream.model.fleet.debezium.DeserializedDebeziumMessage; +import com.openframe.stream.model.fleet.debezium.IntegratedToolEnrichedData; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang3.StringUtils; +import org.springframework.stereotype.Service; + +/** + * Passthrough enrichment for events whose producer already stamped tenant/organization/user + * fields into the message (e.g. Microsoft 365 directory audit events polled per organization). + * Unlike {@link IntegratedToolDataEnrichmentService} it performs no agent/machine lookup — + * such events have no agent, so the deserialized message is the source of truth. + */ +@Slf4j +@Service +@RequiredArgsConstructor +public class PreEnrichedDataEnrichmentService implements DataEnrichmentService { + + private final TenantIdProvider tenantIdProvider; + + @Override + public IntegratedToolEnrichedData getExtraParams(DeserializedDebeziumMessage message) { + IntegratedToolEnrichedData enriched = new IntegratedToolEnrichedData(); + if (message == null) { + return enriched; + } + enriched.setOrganizationId(message.getOrganizationId()); + enriched.setOrganizationName(message.getOrganizationName()); + enriched.setUserId(message.getUserId()); + + String tenantId = StringUtils.isNotBlank(message.getTenantId()) + ? message.getTenantId() + : tenantIdProvider.getTenantId(); + enriched.setTenantId(tenantId); + message.setTenantId(tenantId); + return enriched; + } + + @Override + public DataEnrichmentServiceType getType() { + return DataEnrichmentServiceType.PRE_ENRICHED; + } +} diff --git a/openframe-stream-service-core/src/test/java/com/openframe/stream/deserializer/GoogleWorkspaceAuditEventDeserializerTest.java b/openframe-stream-service-core/src/test/java/com/openframe/stream/deserializer/GoogleWorkspaceAuditEventDeserializerTest.java new file mode 100644 index 000000000..1b16646f0 --- /dev/null +++ b/openframe-stream-service-core/src/test/java/com/openframe/stream/deserializer/GoogleWorkspaceAuditEventDeserializerTest.java @@ -0,0 +1,150 @@ +package com.openframe.stream.deserializer; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.openframe.data.cassandra.model.enums.UnifiedEventType; +import com.openframe.data.model.enums.IntegratedToolType; +import com.openframe.data.model.enums.MessageType; +import com.openframe.kafka.model.debezium.CommonDebeziumMessage; +import com.openframe.kafka.model.debezium.DebeziumMessage; +import com.openframe.stream.model.fleet.debezium.DeserializedDebeziumMessage; +import org.junit.jupiter.api.Test; + +import java.time.Instant; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class GoogleWorkspaceAuditEventDeserializerTest { + + private static final long PROCESSING_TS = 1753868000000L; + + private final ObjectMapper mapper = new ObjectMapper(); + private final GoogleWorkspaceAuditEventDeserializer deserializer = new GoogleWorkspaceAuditEventDeserializer(mapper); + + private static final String AUDIT_EVENT_JSON = """ + { + "uniqueQualifier": "1234567890", "eventIndex": 0, "activityTime": "2026-07-30T10:00:00Z", + "eventType": "USER_SETTINGS", "eventName": "CREATE_USER", "actorEmail": "admin@x.com", + "ipAddress": "203.0.113.5", + "event": {"parameters": [{"name": "USER_EMAIL", "value": "newuser@x.com"}]}, + "tenantId": "tenant-1", "organizationId": "org-uuid-1", "organizationName": "Acme Org", + "connectionId": "conn-1", "connectionName": "Main" + } + """; + + private CommonDebeziumMessage message(String afterJson) { + try { + DebeziumMessage.Payload payload = new DebeziumMessage.Payload<>(); + payload.setAfter(afterJson == null ? null : mapper.readTree(afterJson)); + payload.setOperation("c"); + payload.setTimestamp(PROCESSING_TS); + CommonDebeziumMessage message = new CommonDebeziumMessage(); + message.setPayload(payload); + return message; + } catch (Exception e) { + throw new RuntimeException(e); + } + } + + private DeserializedDebeziumMessage deserialize(String afterJson) { + return deserializer.deserialize(message(afterJson), MessageType.GOOGLE_WORKSPACE_AUDIT_EVENT); + } + + @Test + void registersForGoogleWorkspaceAuditEventType() { + assertEquals(MessageType.GOOGLE_WORKSPACE_AUDIT_EVENT, deserializer.getType()); + } + + @Test + void mapsAuditFieldsToDeserializedMessage() { + DeserializedDebeziumMessage result = deserialize(AUDIT_EVENT_JSON); + + assertEquals("1234567890-0", result.getToolEventId()); + assertEquals(Instant.parse("2026-07-30T10:00:00Z").toEpochMilli(), result.getEventTimestamp()); + assertEquals("2026-07-30", result.getIngestDay()); + assertEquals("CREATE_USER", result.getMessage()); + assertEquals("USER_SETTINGS", result.getSourceEventType()); + assertEquals(UnifiedEventType.GWS_USER_MANAGEMENT, result.getUnifiedEventType()); + assertEquals(IntegratedToolType.GOOGLE_WORKSPACE, result.getIntegratedToolType()); + } + + @Test + void passesThroughTenantAndOrgFields() { + DeserializedDebeziumMessage result = deserialize(AUDIT_EVENT_JSON); + + assertEquals("tenant-1", result.getTenantId()); + assertEquals("org-uuid-1", result.getOrganizationId()); + assertEquals("Acme Org", result.getOrganizationName()); + } + + @Test + void extractsUserIdFromActorEmail() { + DeserializedDebeziumMessage result = deserialize(AUDIT_EVENT_JSON); + + assertEquals("admin@x.com", result.getUserId()); + } + + @Test + void hasNoAgentAndIsVisible() { + DeserializedDebeziumMessage result = deserialize(AUDIT_EVENT_JSON); + + assertNull(result.getAgentId()); + assertTrue(result.getIsVisible()); + assertFalse(result.getSkipProcessing()); + } + + @Test + void detailsContainEventIpAddressAndConnectionFields() throws Exception { + DeserializedDebeziumMessage result = deserialize(AUDIT_EVENT_JSON); + + JsonNode details = mapper.readTree(result.getDetails()); + assertEquals("USER_EMAIL", details.path("event").path("parameters").path(0).path("name").asText()); + assertEquals("203.0.113.5", details.get("ipAddress").asText()); + assertEquals("conn-1", details.get("connectionId").asText()); + assertEquals("Main", details.get("connectionName").asText()); + } + + @Test + void failureEventNameMapsToAuditFailureRegardlessOfEventType() { + String failed = AUDIT_EVENT_JSON.replace("\"eventName\": \"CREATE_USER\"", "\"eventName\": \"LOGIN_FAILURE\""); + + DeserializedDebeziumMessage result = deserialize(failed); + + assertEquals(UnifiedEventType.GWS_AUDIT_FAILURE, result.getUnifiedEventType()); + } + + @Test + void failureMarkerIsCaseInsensitive() { + String failed = AUDIT_EVENT_JSON.replace("\"eventName\": \"CREATE_USER\"", "\"eventName\": \"login_failure\""); + + DeserializedDebeziumMessage result = deserialize(failed); + + assertEquals(UnifiedEventType.GWS_AUDIT_FAILURE, result.getUnifiedEventType()); + } + + @Test + void unmappedEventTypeFallsBackToAuditOther() { + String other = AUDIT_EVENT_JSON.replace("\"eventType\": \"USER_SETTINGS\"", "\"eventType\": \"CALENDAR_SETTINGS\""); + + DeserializedDebeziumMessage result = deserialize(other); + + assertEquals(UnifiedEventType.GWS_AUDIT_OTHER, result.getUnifiedEventType()); + } + + @Test + void missingActivityTimeFallsBackToProcessingTimestamp() { + String withoutTime = AUDIT_EVENT_JSON.replace("\"activityTime\": \"2026-07-30T10:00:00Z\",", ""); + + DeserializedDebeziumMessage result = deserialize(withoutTime); + + assertEquals(PROCESSING_TS, result.getEventTimestamp()); + } + + @Test + void nullAfterReturnsNull() { + assertNull(deserialize(null)); + } +} diff --git a/openframe-stream-service-core/src/test/java/com/openframe/stream/deserializer/Microsoft365AuditEventDeserializerTest.java b/openframe-stream-service-core/src/test/java/com/openframe/stream/deserializer/Microsoft365AuditEventDeserializerTest.java new file mode 100644 index 000000000..390ad39cc --- /dev/null +++ b/openframe-stream-service-core/src/test/java/com/openframe/stream/deserializer/Microsoft365AuditEventDeserializerTest.java @@ -0,0 +1,152 @@ +package com.openframe.stream.deserializer; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.openframe.data.cassandra.model.enums.UnifiedEventType; +import com.openframe.data.model.enums.IntegratedToolType; +import com.openframe.data.model.enums.MessageType; +import com.openframe.kafka.model.debezium.CommonDebeziumMessage; +import com.openframe.kafka.model.debezium.DebeziumMessage; +import com.openframe.stream.model.fleet.debezium.DeserializedDebeziumMessage; +import org.junit.jupiter.api.Test; + +import java.time.Instant; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class Microsoft365AuditEventDeserializerTest { + + private static final long PROCESSING_TS = 1753868000000L; + + private final ObjectMapper mapper = new ObjectMapper(); + private final Microsoft365AuditEventDeserializer deserializer = new Microsoft365AuditEventDeserializer(mapper); + + private static final String AUDIT_EVENT_JSON = """ + { + "auditId": "Directory_abc_123", "activityDateTime": "2026-07-30T10:00:00Z", + "activityDisplayName": "Add user", "category": "UserManagement", "result": "success", + "initiatedBy": {"user": {"userPrincipalName": "admin@x.com", "id": "user-id-1"}}, + "targetResources": [{"displayName": "Test User", "type": "User", "id": "target-id-1"}], + "additionalDetails": [{"key": "UserType", "value": "Member"}], + "tenantId": "tenant-1", "organizationId": "org-uuid-1", "organizationName": "Acme Org", + "connectionId": "conn-1", "connectionName": "Main" + } + """; + + private CommonDebeziumMessage message(String afterJson) { + try { + DebeziumMessage.Payload payload = new DebeziumMessage.Payload<>(); + payload.setAfter(afterJson == null ? null : mapper.readTree(afterJson)); + payload.setOperation("c"); + payload.setTimestamp(PROCESSING_TS); + CommonDebeziumMessage message = new CommonDebeziumMessage(); + message.setPayload(payload); + return message; + } catch (Exception e) { + throw new RuntimeException(e); + } + } + + private DeserializedDebeziumMessage deserialize(String afterJson) { + return deserializer.deserialize(message(afterJson), MessageType.MICROSOFT_365_AUDIT_EVENT); + } + + @Test + void registersForMicrosoft365AuditEventType() { + assertEquals(MessageType.MICROSOFT_365_AUDIT_EVENT, deserializer.getType()); + } + + @Test + void mapsAuditFieldsToDeserializedMessage() { + DeserializedDebeziumMessage result = deserialize(AUDIT_EVENT_JSON); + + assertEquals("Directory_abc_123", result.getToolEventId()); + assertEquals(Instant.parse("2026-07-30T10:00:00Z").toEpochMilli(), result.getEventTimestamp()); + assertEquals("2026-07-30", result.getIngestDay()); + assertEquals("Add user", result.getMessage()); + assertEquals("UserManagement", result.getSourceEventType()); + assertEquals(UnifiedEventType.M365_USER_MANAGEMENT, result.getUnifiedEventType()); + assertEquals(IntegratedToolType.MICROSOFT_365, result.getIntegratedToolType()); + } + + @Test + void passesThroughTenantAndOrgFields() { + DeserializedDebeziumMessage result = deserialize(AUDIT_EVENT_JSON); + + assertEquals("tenant-1", result.getTenantId()); + assertEquals("org-uuid-1", result.getOrganizationId()); + assertEquals("Acme Org", result.getOrganizationName()); + } + + @Test + void extractsUserIdFromInitiatedByUserPrincipalName() { + DeserializedDebeziumMessage result = deserialize(AUDIT_EVENT_JSON); + + assertEquals("admin@x.com", result.getUserId()); + } + + @Test + void hasNoAgentAndIsVisible() { + DeserializedDebeziumMessage result = deserialize(AUDIT_EVENT_JSON); + + assertNull(result.getAgentId()); + assertTrue(result.getIsVisible()); + assertFalse(result.getSkipProcessing()); + } + + @Test + void detailsContainInitiatedByTargetResourcesAndAdditionalDetails() throws Exception { + DeserializedDebeziumMessage result = deserialize(AUDIT_EVENT_JSON); + + JsonNode details = mapper.readTree(result.getDetails()); + assertEquals("admin@x.com", details.path("initiatedBy").path("user").path("userPrincipalName").asText()); + assertEquals("User", details.path("targetResources").path(0).path("type").asText()); + assertEquals("UserType", details.path("additionalDetails").path(0).path("key").asText()); + assertEquals("conn-1", details.get("connectionId").asText()); + assertEquals("Main", details.get("connectionName").asText()); + } + + @Test + void failureResultMapsToAuditFailureRegardlessOfCategory() { + String failed = AUDIT_EVENT_JSON.replace("\"result\": \"success\"", "\"result\": \"failure\""); + + DeserializedDebeziumMessage result = deserialize(failed); + + assertEquals(UnifiedEventType.M365_AUDIT_FAILURE, result.getUnifiedEventType()); + } + + @Test + void timeoutResultMapsToAuditFailure() { + String timedOut = AUDIT_EVENT_JSON.replace("\"result\": \"success\"", "\"result\": \"timeout\""); + + DeserializedDebeziumMessage result = deserialize(timedOut); + + assertEquals(UnifiedEventType.M365_AUDIT_FAILURE, result.getUnifiedEventType()); + } + + @Test + void unmappedCategoryFallsBackToAuditOther() { + String other = AUDIT_EVENT_JSON.replace("\"category\": \"UserManagement\"", "\"category\": \"KeyManagement\""); + + DeserializedDebeziumMessage result = deserialize(other); + + assertEquals(UnifiedEventType.M365_AUDIT_OTHER, result.getUnifiedEventType()); + } + + @Test + void missingActivityDateTimeFallsBackToProcessingTimestamp() { + String withoutTime = AUDIT_EVENT_JSON.replace("\"activityDateTime\": \"2026-07-30T10:00:00Z\",", ""); + + DeserializedDebeziumMessage result = deserialize(withoutTime); + + assertEquals(PROCESSING_TS, result.getEventTimestamp()); + } + + @Test + void nullAfterReturnsNull() { + assertNull(deserialize(null)); + } +} diff --git a/openframe-stream-service-core/src/test/java/com/openframe/stream/mapping/EventTypeMapperTest.java b/openframe-stream-service-core/src/test/java/com/openframe/stream/mapping/EventTypeMapperTest.java new file mode 100644 index 000000000..782645450 --- /dev/null +++ b/openframe-stream-service-core/src/test/java/com/openframe/stream/mapping/EventTypeMapperTest.java @@ -0,0 +1,57 @@ +package com.openframe.stream.mapping; + +import com.openframe.data.cassandra.model.enums.UnifiedEventType; +import com.openframe.data.model.enums.IntegratedToolType; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.CsvSource; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +class EventTypeMapperTest { + + @ParameterizedTest + @CsvSource({ + "UserManagement, M365_USER_MANAGEMENT", + "GroupManagement, M365_GROUP_MANAGEMENT", + "ApplicationManagement, M365_APPLICATION_MANAGEMENT", + "RoleManagement, M365_ROLE_MANAGEMENT", + "Policy, M365_POLICY", + "DirectoryManagement, M365_DIRECTORY_MANAGEMENT" + }) + void mapsMicrosoft365CategoriesToUnifiedTypes(String category, UnifiedEventType expected) { + assertEquals(expected, EventTypeMapper.mapToUnifiedType(IntegratedToolType.MICROSOFT_365, category)); + } + + @Test + void unmappedMicrosoft365CategoryFallsBackToUnknown() { + assertEquals(UnifiedEventType.UNKNOWN, + EventTypeMapper.mapToUnifiedType(IntegratedToolType.MICROSOFT_365, "KeyManagement")); + } + + @Test + void categoryMappingIsToolScoped() { + assertEquals(UnifiedEventType.UNKNOWN, + EventTypeMapper.mapToUnifiedType(IntegratedToolType.FLEET, "UserManagement")); + } + + @ParameterizedTest + @CsvSource({ + "USER_SETTINGS, GWS_USER_MANAGEMENT", + "GROUP_SETTINGS, GWS_GROUP_MANAGEMENT", + "SECURITY_SETTINGS, GWS_SECURITY_SETTINGS", + "DOMAIN_SETTINGS, GWS_DOMAIN_SETTINGS", + "DELEGATED_ADMIN_SETTINGS, GWS_DELEGATED_ADMIN" + }) + void mapsGoogleWorkspaceEventTypesToUnifiedTypes(String eventType, UnifiedEventType expected) { + assertEquals(expected, EventTypeMapper.mapToUnifiedType(IntegratedToolType.GOOGLE_WORKSPACE, eventType)); + } + + @Test + void googleWorkspaceMappingIsToolScoped() { + assertEquals(UnifiedEventType.UNKNOWN, + EventTypeMapper.mapToUnifiedType(IntegratedToolType.GOOGLE_WORKSPACE, "UserManagement")); + assertEquals(UnifiedEventType.UNKNOWN, + EventTypeMapper.mapToUnifiedType(IntegratedToolType.MICROSOFT_365, "USER_SETTINGS")); + } +} diff --git a/openframe-stream-service-core/src/test/java/com/openframe/stream/service/PreEnrichedDataEnrichmentServiceTest.java b/openframe-stream-service-core/src/test/java/com/openframe/stream/service/PreEnrichedDataEnrichmentServiceTest.java new file mode 100644 index 000000000..c20258cfd --- /dev/null +++ b/openframe-stream-service-core/src/test/java/com/openframe/stream/service/PreEnrichedDataEnrichmentServiceTest.java @@ -0,0 +1,68 @@ +package com.openframe.stream.service; + +import com.openframe.data.model.enums.DataEnrichmentServiceType; +import com.openframe.data.service.TenantIdProvider; +import com.openframe.stream.model.fleet.debezium.DeserializedDebeziumMessage; +import com.openframe.stream.model.fleet.debezium.IntegratedToolEnrichedData; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +class PreEnrichedDataEnrichmentServiceTest { + + private TenantIdProvider tenantIdProvider; + private PreEnrichedDataEnrichmentService service; + + @BeforeEach + void setUp() { + tenantIdProvider = mock(TenantIdProvider.class); + service = new PreEnrichedDataEnrichmentService(tenantIdProvider); + } + + @Test + void registersAsPreEnrichedType() { + assertEquals(DataEnrichmentServiceType.PRE_ENRICHED, service.getType()); + } + + @Test + void passesThroughPreEnrichedFieldsFromMessage() { + DeserializedDebeziumMessage message = DeserializedDebeziumMessage.builder() + .tenantId("tenant-1") + .organizationId("org-uuid-1") + .organizationName("Acme Org") + .userId("admin@x.com") + .build(); + + IntegratedToolEnrichedData enriched = service.getExtraParams(message); + + assertEquals("tenant-1", enriched.getTenantId()); + assertEquals("org-uuid-1", enriched.getOrganizationId()); + assertEquals("Acme Org", enriched.getOrganizationName()); + assertEquals("admin@x.com", enriched.getUserId()); + } + + @Test + void missingTenantFallsBackToTenantIdProvider() { + when(tenantIdProvider.getTenantId()).thenReturn("provider-tenant"); + DeserializedDebeziumMessage message = DeserializedDebeziumMessage.builder() + .organizationId("org-uuid-1") + .build(); + + IntegratedToolEnrichedData enriched = service.getExtraParams(message); + + assertEquals("provider-tenant", enriched.getTenantId()); + assertEquals("provider-tenant", message.getTenantId()); + } + + @Test + void nullMessageYieldsEmptyEnrichment() { + IntegratedToolEnrichedData enriched = service.getExtraParams(null); + + assertNull(enriched.getTenantId()); + assertNull(enriched.getOrganizationId()); + } +}