From 6d1fe68b4da82100bc704509da82eb520d7a162a Mon Sep 17 00:00:00 2001 From: andrii Date: Thu, 30 Jul 2026 14:09:20 +0300 Subject: [PATCH 1/9] feat(stream): microsoft 365 directory-audit event type and deserializer Additive support for Entra directory audit events entering the logs pipeline via the generic message-type path: - IntegratedToolType.MICROSOFT_365 ("microsoft-365") - MessageType.MICROSOFT_365_AUDIT_EVENT -> CASSANDRA_EVENT_LOG + KAFKA_PINOT - UnifiedEventType M365_* additions (RoleManagement/failures -> WARNING) - EventTypeMapper mappings for Graph directoryAudits category values - Microsoft365AuditEventDeserializer: toolEventId = Graph audit id (idempotent upserts), eventTimestamp from activityDateTime, result=failure -> M365_AUDIT_FAILURE, unmapped category -> M365_AUDIT_OTHER A2.2 enrichment decision: INTEGRATED_TOOLS_EVENTS cannot supply org fields for agentless events (machine-lookup only populates enriched data when agentId is present), so a new PRE_ENRICHED DataEnrichmentServiceType + PreEnrichedDataEnrichmentService passes organizationId/organizationName/ userId through from the deserialized message (populated from the pre-enriched payload), with TenantIdProvider fallback for tenantId. DeserializedDebeziumMessage gains organizationId/organizationName/userId (additive; existing MessageTypes unaffected). Co-Authored-By: Claude Fable 5 --- .../model/enums/UnifiedEventType.java | 10 ++ .../enums/DataEnrichmentServiceType.java | 3 +- .../data/model/enums/IntegratedToolType.java | 3 +- .../data/model/enums/MessageType.java | 2 + .../Microsoft365AuditEventDeserializer.java | 121 +++++++++++++++ .../stream/mapping/EventTypeMapper.java | 8 + .../stream/mapping/SourceEventTypes.java | 13 ++ .../debezium/DeserializedDebeziumMessage.java | 3 + .../PreEnrichedDataEnrichmentService.java | 47 ++++++ ...icrosoft365AuditEventDeserializerTest.java | 138 ++++++++++++++++++ .../stream/mapping/EventTypeMapperTest.java | 37 +++++ .../PreEnrichedDataEnrichmentServiceTest.java | 68 +++++++++ 12 files changed, 451 insertions(+), 2 deletions(-) create mode 100644 openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/Microsoft365AuditEventDeserializer.java create mode 100644 openframe-stream-service-core/src/main/java/com/openframe/stream/service/PreEnrichedDataEnrichmentService.java create mode 100644 openframe-stream-service-core/src/test/java/com/openframe/stream/deserializer/Microsoft365AuditEventDeserializerTest.java create mode 100644 openframe-stream-service-core/src/test/java/com/openframe/stream/mapping/EventTypeMapperTest.java create mode 100644 openframe-stream-service-core/src/test/java/com/openframe/stream/service/PreEnrichedDataEnrichmentServiceTest.java 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..c549f7c32 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,16 @@ 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"), + // 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..a6820527a 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,8 @@ public enum IntegratedToolType { RMM("rmm"), MESHCENTRAL ("meshcentral"), - FLEET ("fleet-mdm"); + FLEET ("fleet-mdm"), + MICROSOFT_365("microsoft-365"); 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..d041190ee 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,8 @@ 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); private final IntegratedToolType integratedToolType; 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..fc63a4b38 --- /dev/null +++ b/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/Microsoft365AuditEventDeserializer.java @@ -0,0 +1,121 @@ +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 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. + */ +@Slf4j +@Component +@RequiredArgsConstructor +public class Microsoft365AuditEventDeserializer implements KafkaMessageDeserializer { + + private static final DateTimeFormatter DAY_FORMATTER = + DateTimeFormatter.ofPattern("yyyy-MM-dd").withZone(ZoneId.of("UTC")); + private static final String RESULT_FAILURE = "failure"; + 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").filter(RESULT_FAILURE::equalsIgnoreCase).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); + } + 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..1b5e0de52 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,13 @@ 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); } } 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..66e321eee 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,17 @@ 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"; + } } 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/Microsoft365AuditEventDeserializerTest.java b/openframe-stream-service-core/src/test/java/com/openframe/stream/deserializer/Microsoft365AuditEventDeserializerTest.java new file mode 100644 index 000000000..590bbc1ec --- /dev/null +++ b/openframe-stream-service-core/src/test/java/com/openframe/stream/deserializer/Microsoft365AuditEventDeserializerTest.java @@ -0,0 +1,138 @@ +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"}], + "tenantId": "tenant-1", "organizationId": "org-uuid-1", "organizationName": "Acme Org" + } + """; + + 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 detailsContainInitiatedByAndTargetResources() 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()); + } + + @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 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..097995a94 --- /dev/null +++ b/openframe-stream-service-core/src/test/java/com/openframe/stream/mapping/EventTypeMapperTest.java @@ -0,0 +1,37 @@ +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")); + } +} 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()); + } +} From d753f279d56b133fcb57fd6060cd10435d37921e Mon Sep 17 00:00:00 2001 From: andrii Date: Thu, 30 Jul 2026 14:54:04 +0300 Subject: [PATCH 2/9] fix(stream): classify timeout audit results as failures, preserve additionalDetails MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review follow-up (PR #1603): - Graph directoryAudits result can be success|failure|timeout — timeout now also maps to M365_AUDIT_FAILURE instead of falling through to the category mapping as an INFO event - details JSON now carries provider-supplied additionalDetails alongside initiatedBy and targetResources (payload contract extended in the poller and Graph model accordingly) Co-Authored-By: Claude Fable 5 --- .../Microsoft365AuditEventDeserializer.java | 11 +++++++++-- .../Microsoft365AuditEventDeserializerTest.java | 13 ++++++++++++- 2 files changed, 21 insertions(+), 3 deletions(-) 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 index fc63a4b38..a2995f8c5 100644 --- 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 @@ -17,6 +17,7 @@ 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 @@ -33,7 +34,9 @@ public class Microsoft365AuditEventDeserializer implements KafkaMessageDeseriali private static final DateTimeFormatter DAY_FORMATTER = DateTimeFormatter.ofPattern("yyyy-MM-dd").withZone(ZoneId.of("UTC")); - private static final String RESULT_FAILURE = "failure"; + // 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; @@ -75,7 +78,7 @@ public DeserializedDebeziumMessage deserialize(CommonDebeziumMessage debeziumMes } private UnifiedEventType resolveEventType(JsonNode after, String category) { - if (textField(after, "result").filter(RESULT_FAILURE::equalsIgnoreCase).isPresent()) { + if (textField(after, "result").map(String::toLowerCase).filter(FAILED_RESULTS::contains).isPresent()) { return UnifiedEventType.M365_AUDIT_FAILURE; } UnifiedEventType mapped = EventTypeMapper.mapToUnifiedType(getType().getIntegratedToolType(), category); @@ -104,6 +107,10 @@ private String buildDetails(JsonNode after) { if (targetResources != null && !targetResources.isNull()) { details.set("targetResources", targetResources); } + JsonNode additionalDetails = after.get("additionalDetails"); + if (additionalDetails != null && !additionalDetails.isNull()) { + details.set("additionalDetails", additionalDetails); + } return details.toString(); } 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 index 590bbc1ec..1c88934a8 100644 --- 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 @@ -30,6 +30,7 @@ class Microsoft365AuditEventDeserializerTest { "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" } """; @@ -96,12 +97,13 @@ void hasNoAgentAndIsVisible() { } @Test - void detailsContainInitiatedByAndTargetResources() throws Exception { + 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()); } @Test @@ -113,6 +115,15 @@ void failureResultMapsToAuditFailureRegardlessOfCategory() { 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\""); From c2fb18ebb5044b0d8894614bb2b2dc89c97852d9 Mon Sep 17 00:00:00 2001 From: andrii Date: Mon, 3 Aug 2026 15:44:48 +0300 Subject: [PATCH 3/9] feat: pass Microsoft 365 audit connection identity through to log details Co-Authored-By: Claude Fable 5 --- .../deserializer/Microsoft365AuditEventDeserializer.java | 3 +++ .../deserializer/Microsoft365AuditEventDeserializerTest.java | 5 ++++- 2 files changed, 7 insertions(+), 1 deletion(-) 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 index a2995f8c5..2e478de1d 100644 --- 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 @@ -26,6 +26,7 @@ * 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 @@ -111,6 +112,8 @@ private String buildDetails(JsonNode after) { 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(); } 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 index 1c88934a8..390ad39cc 100644 --- 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 @@ -31,7 +31,8 @@ class Microsoft365AuditEventDeserializerTest { "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" + "tenantId": "tenant-1", "organizationId": "org-uuid-1", "organizationName": "Acme Org", + "connectionId": "conn-1", "connectionName": "Main" } """; @@ -104,6 +105,8 @@ void detailsContainInitiatedByTargetResourcesAndAdditionalDetails() throws Excep 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 From 6ad81c36472e28de9602a69e01f826ffeb14c29f Mon Sep 17 00:00:00 2001 From: andrii Date: Mon, 10 Aug 2026 14:58:29 +0200 Subject: [PATCH 4/9] feat(kafka,cassandra): Google Workspace audit event enums Co-Authored-By: Claude Fable 5 --- .../data/cassandra/model/enums/UnifiedEventType.java | 9 +++++++++ .../openframe/data/model/enums/IntegratedToolType.java | 3 ++- .../java/com/openframe/data/model/enums/MessageType.java | 2 ++ 3 files changed, 13 insertions(+), 1 deletion(-) 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 c549f7c32..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 @@ -174,6 +174,15 @@ public enum UnifiedEventType { 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/IntegratedToolType.java b/openframe-data-kafka/src/main/java/com/openframe/data/model/enums/IntegratedToolType.java index a6820527a..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 @@ -5,7 +5,8 @@ public enum IntegratedToolType { RMM("rmm"), MESHCENTRAL ("meshcentral"), FLEET ("fleet-mdm"), - MICROSOFT_365("microsoft-365"); + 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 d041190ee..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 @@ -22,6 +22,8 @@ public enum MessageType { 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; From 4be1ac3f084c2dda9bad5b66585a75f1d11e1e59 Mon Sep 17 00:00:00 2001 From: andrii Date: Mon, 10 Aug 2026 15:02:25 +0200 Subject: [PATCH 5/9] feat(stream): Google Workspace source event types + mappings Co-Authored-By: Claude Fable 5 --- .../stream/mapping/EventTypeMapper.java | 7 +++++++ .../stream/mapping/SourceEventTypes.java | 12 +++++++++++ .../stream/mapping/EventTypeMapperTest.java | 20 +++++++++++++++++++ 3 files changed, 39 insertions(+) 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 1b5e0de52..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 @@ -277,5 +277,12 @@ private static void initializeDefaultMappings() { 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 66e321eee..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 @@ -263,4 +263,16 @@ interface Microsoft365 { 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/test/java/com/openframe/stream/mapping/EventTypeMapperTest.java b/openframe-stream-service-core/src/test/java/com/openframe/stream/mapping/EventTypeMapperTest.java index 097995a94..782645450 100644 --- 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 @@ -34,4 +34,24 @@ 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")); + } } From 708eccbbce9e42555c95b1f5377078179750a715 Mon Sep 17 00:00:00 2001 From: andrii Date: Mon, 10 Aug 2026 15:10:18 +0200 Subject: [PATCH 6/9] feat(stream): Google Workspace directory-audit deserializer Add GoogleWorkspaceAuditEventDeserializer mirroring the Microsoft 365 Entra audit deserializer: events are hand-built by the poller (not CDC), pre-enriched with tenant/organization fields, and carry no agent reference. KafkaMessageDeserializer is single-result, so this stays 1:1 with the Kafka record; fanning one polled activity's events[] into N records is the poller's job (Phase C). toolEventId = uniqueQualifier + "-" + eventIndex. Failure predicate: eventName containing "_FAILURE" (case-insensitive) -> GWS_AUDIT_FAILURE before category mapping; unmapped eventType -> GWS_AUDIT_OTHER, never UNKNOWN. Ports the MS365 deserializer's 11 test invariants to the Google field map. Also backfills the .Microsoft365AuditEventDeserializer.md sidecar doc that was missing, per module convention. Co-Authored-By: Claude Fable 5 --- .../.GoogleWorkspaceAuditEventDeserializer.md | 71 +++++++++ .../.Microsoft365AuditEventDeserializer.md | 53 +++++++ ...GoogleWorkspaceAuditEventDeserializer.java | 134 ++++++++++++++++ ...leWorkspaceAuditEventDeserializerTest.java | 150 ++++++++++++++++++ 4 files changed, 408 insertions(+) create mode 100644 openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/.GoogleWorkspaceAuditEventDeserializer.md create mode 100644 openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/.Microsoft365AuditEventDeserializer.md create mode 100644 openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/GoogleWorkspaceAuditEventDeserializer.java create mode 100644 openframe-stream-service-core/src/test/java/com/openframe/stream/deserializer/GoogleWorkspaceAuditEventDeserializerTest.java diff --git a/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/.GoogleWorkspaceAuditEventDeserializer.md b/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/.GoogleWorkspaceAuditEventDeserializer.md new file mode 100644 index 000000000..13dcb49ca --- /dev/null +++ b/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/.GoogleWorkspaceAuditEventDeserializer.md @@ -0,0 +1,71 @@ + +A Spring component that deserializes Google Workspace directory audit events polled from the Admin SDK Reports API `activities.list("admin")` endpoint. Mirrors `Microsoft365AuditEventDeserializer`'s pattern: events are hand-built by the poller (not Debezium CDC), arrive pre-enriched with tenant/organization fields, and carry no agent reference — routed via `DataEnrichmentServiceType.PRE_ENRICHED`. + +## Poller Contract + +`KafkaMessageDeserializer.deserialize()` is single-result, so this deserializer stays 1:1 with the incoming Kafka record — the fan-out from one Reports API activity's `events[]` array into N Kafka records is the poller's responsibility, not this class's. The poller writes a FLAT `after` object per event with: + +| `after` field | Meaning | +|---|---| +| `uniqueQualifier` | Reports API activity identifier | +| `eventIndex` | Index of this event within the activity's `events[]` array | +| `activityTime` | RFC3339 timestamp of the activity | +| `eventType` | Reports API `events[].type` (e.g. `USER_SETTINGS`, `GROUP_SETTINGS`) | +| `eventName` | Reports API `events[].name` (e.g. `CREATE_USER`, `LOGIN_FAILURE`) | +| `actorEmail` | Acting admin/user email | +| `ipAddress` | Source IP of the action | +| `event` | Nested JSON carrying the raw event's `parameters` | +| `tenantId` / `organizationId` / `organizationName` | Tenant/org passthrough | +| `connectionId` / `connectionName` | Multi-connection org passthrough | + +## Key Components + +- **`getType()`** — Returns `MessageType.GOOGLE_WORKSPACE_AUDIT_EVENT` +- **`deserialize()`** — Maps the flat `after` payload built by the poller into a `DeserializedDebeziumMessage`; returns `null` when `after` is null +- **`resolveEventType()`** — Failure-first classification: `eventName` containing `_FAILURE` (case-insensitive) always maps to `UnifiedEventType.GWS_AUDIT_FAILURE` regardless of `eventType`; otherwise delegates to `EventTypeMapper` and falls back to `GWS_AUDIT_OTHER` for unmapped event types (never `UNKNOWN`) +- **`buildToolEventId()`** — `uniqueQualifier + "-" + eventIndex`; a single activity can carry multiple events, so the pair keeps replays from the poller's cursor overlap window upsert idempotent per event +- **`getEventTimestamp()`** — Parses `activityTime` (RFC3339/ISO-8601); falls back to the Debezium payload's processing `timestamp` when absent or unparseable +- **`buildDetails()`** — Verbatim `event` sub-object (parameters), plus flat `ipAddress`, `connectionId`/`connectionName` (multi-connection orgs) + +## Field Map + +| `DeserializedDebeziumMessage` field | Source in `after` | +|---|---| +| `toolEventId` | `uniqueQualifier + "-" + eventIndex` | +| `eventTimestamp` | `activityTime`, fallback processing timestamp | +| `message` | `eventName` | +| `sourceEventType` | `eventType` | +| `userId` | `actorEmail` | +| `tenantId` / `organizationId` / `organizationName` | passthrough | +| `agentId` | always `null` | +| `isVisible` / `skipProcessing` | always `true` / `false` | + +## Usage Example + +```java +ObjectMapper mapper = new ObjectMapper(); +GoogleWorkspaceAuditEventDeserializer deserializer = new GoogleWorkspaceAuditEventDeserializer(mapper); + +// Spring wires this as the strategy registered for MessageType.GOOGLE_WORKSPACE_AUDIT_EVENT +MessageType type = deserializer.getType(); +// → MessageType.GOOGLE_WORKSPACE_AUDIT_EVENT + +JsonNode after = mapper.readTree(""" + { + "uniqueQualifier": "1234567890", "eventIndex": 0, "activityTime": "2026-07-30T10:00:00Z", + "eventType": "USER_SETTINGS", "eventName": "CREATE_USER", "actorEmail": "admin@x.com", + "event": {"parameters": [{"name": "USER_EMAIL", "value": "newuser@x.com"}]}, + "tenantId": "tenant-1", "organizationId": "org-uuid-1", "organizationName": "Acme Org" + } + """); + +DeserializedDebeziumMessage result = deserializer.deserialize(debeziumMessageWith(after), type); +// → unifiedEventType = GWS_USER_MANAGEMENT, userId = "admin@x.com", toolEventId = "1234567890-0" +``` + +## Related Files + +- [`Microsoft365AuditEventDeserializer`](https://github.com/flamingo-stack/openframe-oss-lib/blob/main/src/main/java/com/openframe/stream/deserializer/Microsoft365AuditEventDeserializer.java) — template this deserializer mirrors +- [`EventTypeMapper`](https://github.com/flamingo-stack/openframe-oss-lib/blob/main/src/main/java/com/openframe/stream/mapping/EventTypeMapper.java) — `eventType` → `UnifiedEventType` registry +- [`SourceEventTypes.GoogleWorkspace`](https://github.com/flamingo-stack/openframe-oss-lib/blob/main/src/main/java/com/openframe/stream/mapping/SourceEventTypes.java) — `eventType` string constants +- [`MessageType`](https://github.com/flamingo-stack/openframe-oss-lib/blob/main/src/main/java/com/openframe/data/model/enums/MessageType.java) — enum of all routable message types diff --git a/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/.Microsoft365AuditEventDeserializer.md b/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/.Microsoft365AuditEventDeserializer.md new file mode 100644 index 000000000..080c5fafe --- /dev/null +++ b/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/.Microsoft365AuditEventDeserializer.md @@ -0,0 +1,53 @@ + +A Spring component that deserializes Microsoft 365 Entra directory audit events polled from Graph `auditLogs/directoryAudits`. Unlike the CDC-sourced deserializers in this package, the events are hand-built by the poller (not Debezium change capture), arrive pre-enriched with tenant/organization fields, and carry no agent reference — routed via `DataEnrichmentServiceType.PRE_ENRICHED`. + +## Key Components + +- **`getType()`** — Returns `MessageType.MICROSOFT_365_AUDIT_EVENT` +- **`deserialize()`** — Maps the flat `after` payload built by the poller into a `DeserializedDebeziumMessage`; returns `null` when `after` is null +- **`resolveEventType()`** — Failure-first classification: a `result` of `failure`/`timeout` (case-insensitive) always maps to `UnifiedEventType.M365_AUDIT_FAILURE` regardless of `category`; otherwise delegates to `EventTypeMapper` and falls back to `M365_AUDIT_OTHER` for unmapped categories (never `UNKNOWN`) +- **`getEventTimestamp()`** — Parses `activityDateTime` (RFC3339/ISO-8601); falls back to the Debezium payload's processing `timestamp` when absent or unparseable +- **`buildDetails()`** — Verbatim `initiatedBy`, `targetResources`, `additionalDetails` sub-objects, plus flat `connectionId`/`connectionName` (multi-connection orgs) + +## Field Map + +| `DeserializedDebeziumMessage` field | Source in `after` | +|---|---| +| `toolEventId` | `auditId` (Graph audit record id — idempotent across poller replay) | +| `eventTimestamp` | `activityDateTime`, fallback processing timestamp | +| `message` | `activityDisplayName` | +| `sourceEventType` | `category` | +| `userId` | `initiatedBy.user.userPrincipalName` | +| `tenantId` / `organizationId` / `organizationName` | passthrough | +| `agentId` | always `null` | +| `isVisible` / `skipProcessing` | always `true` / `false` | + +## Usage Example + +```java +ObjectMapper mapper = new ObjectMapper(); +Microsoft365AuditEventDeserializer deserializer = new Microsoft365AuditEventDeserializer(mapper); + +// Spring wires this as the strategy registered for MessageType.MICROSOFT_365_AUDIT_EVENT +MessageType type = deserializer.getType(); +// → MessageType.MICROSOFT_365_AUDIT_EVENT + +JsonNode after = mapper.readTree(""" + { + "auditId": "Directory_abc_123", "activityDateTime": "2026-07-30T10:00:00Z", + "activityDisplayName": "Add user", "category": "UserManagement", "result": "success", + "initiatedBy": {"user": {"userPrincipalName": "admin@x.com"}}, + "tenantId": "tenant-1", "organizationId": "org-uuid-1", "organizationName": "Acme Org" + } + """); + +DeserializedDebeziumMessage result = deserializer.deserialize(debeziumMessageWith(after), type); +// → unifiedEventType = M365_USER_MANAGEMENT, userId = "admin@x.com", toolEventId = "Directory_abc_123" +``` + +## Related Files + +- [`GoogleWorkspaceAuditEventDeserializer`](https://github.com/flamingo-stack/openframe-oss-lib/blob/main/src/main/java/com/openframe/stream/deserializer/GoogleWorkspaceAuditEventDeserializer.java) — sibling deserializer for Google Workspace directory audit events, same poller-fed / pre-enriched / 1:1 pattern +- [`EventTypeMapper`](https://github.com/flamingo-stack/openframe-oss-lib/blob/main/src/main/java/com/openframe/stream/mapping/EventTypeMapper.java) — `category` → `UnifiedEventType` registry +- [`SourceEventTypes.Microsoft365`](https://github.com/flamingo-stack/openframe-oss-lib/blob/main/src/main/java/com/openframe/stream/mapping/SourceEventTypes.java) — `category` string constants +- [`MessageType`](https://github.com/flamingo-stack/openframe-oss-lib/blob/main/src/main/java/com/openframe/data/model/enums/MessageType.java) — enum of all routable message types 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..98aa1e058 --- /dev/null +++ b/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/GoogleWorkspaceAuditEventDeserializer.java @@ -0,0 +1,134 @@ +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. + */ +@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/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)); + } +} From ef5c6ce6cb68da81f98619778582ecd51f5adadc Mon Sep 17 00:00:00 2001 From: andrii Date: Tue, 11 Aug 2026 12:08:27 +0200 Subject: [PATCH 7/9] docs(stream): document the Reports API parameter union in the GW audit deserializer contract Co-Authored-By: Claude Fable 5 --- .../deserializer/.GoogleWorkspaceAuditEventDeserializer.md | 2 +- .../GoogleWorkspaceAuditEventDeserializer.java | 7 +++++++ 2 files changed, 8 insertions(+), 1 deletion(-) diff --git a/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/.GoogleWorkspaceAuditEventDeserializer.md b/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/.GoogleWorkspaceAuditEventDeserializer.md index 13dcb49ca..5c364eb4b 100644 --- a/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/.GoogleWorkspaceAuditEventDeserializer.md +++ b/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/.GoogleWorkspaceAuditEventDeserializer.md @@ -1,4 +1,4 @@ - + A Spring component that deserializes Google Workspace directory audit events polled from the Admin SDK Reports API `activities.list("admin")` endpoint. Mirrors `Microsoft365AuditEventDeserializer`'s pattern: events are hand-built by the poller (not Debezium CDC), arrive pre-enriched with tenant/organization fields, and carry no agent reference — routed via `DataEnrichmentServiceType.PRE_ENRICHED`. ## Poller Contract 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 index 98aa1e058..9afb5f9db 100644 --- 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 @@ -35,6 +35,13 @@ * 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} 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}; nothing may assume {@code value} alone. */ @Slf4j @Component From 3f05d6b7b976851482daa3134ceb3139111208a1 Mon Sep 17 00:00:00 2001 From: andrii Date: Wed, 12 Aug 2026 11:41:52 +0200 Subject: [PATCH 8/9] chore(stream): drop autogenerated deserializer docs from PR These .md contracts next to the deserializers are produced by the automated documentation job, so they should not be hand-carried in a feature PR. Co-Authored-By: Claude Opus 5 --- .../.GoogleWorkspaceAuditEventDeserializer.md | 71 ------------------- .../.Microsoft365AuditEventDeserializer.md | 53 -------------- 2 files changed, 124 deletions(-) delete mode 100644 openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/.GoogleWorkspaceAuditEventDeserializer.md delete mode 100644 openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/.Microsoft365AuditEventDeserializer.md diff --git a/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/.GoogleWorkspaceAuditEventDeserializer.md b/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/.GoogleWorkspaceAuditEventDeserializer.md deleted file mode 100644 index 5c364eb4b..000000000 --- a/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/.GoogleWorkspaceAuditEventDeserializer.md +++ /dev/null @@ -1,71 +0,0 @@ - -A Spring component that deserializes Google Workspace directory audit events polled from the Admin SDK Reports API `activities.list("admin")` endpoint. Mirrors `Microsoft365AuditEventDeserializer`'s pattern: events are hand-built by the poller (not Debezium CDC), arrive pre-enriched with tenant/organization fields, and carry no agent reference — routed via `DataEnrichmentServiceType.PRE_ENRICHED`. - -## Poller Contract - -`KafkaMessageDeserializer.deserialize()` is single-result, so this deserializer stays 1:1 with the incoming Kafka record — the fan-out from one Reports API activity's `events[]` array into N Kafka records is the poller's responsibility, not this class's. The poller writes a FLAT `after` object per event with: - -| `after` field | Meaning | -|---|---| -| `uniqueQualifier` | Reports API activity identifier | -| `eventIndex` | Index of this event within the activity's `events[]` array | -| `activityTime` | RFC3339 timestamp of the activity | -| `eventType` | Reports API `events[].type` (e.g. `USER_SETTINGS`, `GROUP_SETTINGS`) | -| `eventName` | Reports API `events[].name` (e.g. `CREATE_USER`, `LOGIN_FAILURE`) | -| `actorEmail` | Acting admin/user email | -| `ipAddress` | Source IP of the action | -| `event` | Nested JSON carrying the raw event's `parameters` | -| `tenantId` / `organizationId` / `organizationName` | Tenant/org passthrough | -| `connectionId` / `connectionName` | Multi-connection org passthrough | - -## Key Components - -- **`getType()`** — Returns `MessageType.GOOGLE_WORKSPACE_AUDIT_EVENT` -- **`deserialize()`** — Maps the flat `after` payload built by the poller into a `DeserializedDebeziumMessage`; returns `null` when `after` is null -- **`resolveEventType()`** — Failure-first classification: `eventName` containing `_FAILURE` (case-insensitive) always maps to `UnifiedEventType.GWS_AUDIT_FAILURE` regardless of `eventType`; otherwise delegates to `EventTypeMapper` and falls back to `GWS_AUDIT_OTHER` for unmapped event types (never `UNKNOWN`) -- **`buildToolEventId()`** — `uniqueQualifier + "-" + eventIndex`; a single activity can carry multiple events, so the pair keeps replays from the poller's cursor overlap window upsert idempotent per event -- **`getEventTimestamp()`** — Parses `activityTime` (RFC3339/ISO-8601); falls back to the Debezium payload's processing `timestamp` when absent or unparseable -- **`buildDetails()`** — Verbatim `event` sub-object (parameters), plus flat `ipAddress`, `connectionId`/`connectionName` (multi-connection orgs) - -## Field Map - -| `DeserializedDebeziumMessage` field | Source in `after` | -|---|---| -| `toolEventId` | `uniqueQualifier + "-" + eventIndex` | -| `eventTimestamp` | `activityTime`, fallback processing timestamp | -| `message` | `eventName` | -| `sourceEventType` | `eventType` | -| `userId` | `actorEmail` | -| `tenantId` / `organizationId` / `organizationName` | passthrough | -| `agentId` | always `null` | -| `isVisible` / `skipProcessing` | always `true` / `false` | - -## Usage Example - -```java -ObjectMapper mapper = new ObjectMapper(); -GoogleWorkspaceAuditEventDeserializer deserializer = new GoogleWorkspaceAuditEventDeserializer(mapper); - -// Spring wires this as the strategy registered for MessageType.GOOGLE_WORKSPACE_AUDIT_EVENT -MessageType type = deserializer.getType(); -// → MessageType.GOOGLE_WORKSPACE_AUDIT_EVENT - -JsonNode after = mapper.readTree(""" - { - "uniqueQualifier": "1234567890", "eventIndex": 0, "activityTime": "2026-07-30T10:00:00Z", - "eventType": "USER_SETTINGS", "eventName": "CREATE_USER", "actorEmail": "admin@x.com", - "event": {"parameters": [{"name": "USER_EMAIL", "value": "newuser@x.com"}]}, - "tenantId": "tenant-1", "organizationId": "org-uuid-1", "organizationName": "Acme Org" - } - """); - -DeserializedDebeziumMessage result = deserializer.deserialize(debeziumMessageWith(after), type); -// → unifiedEventType = GWS_USER_MANAGEMENT, userId = "admin@x.com", toolEventId = "1234567890-0" -``` - -## Related Files - -- [`Microsoft365AuditEventDeserializer`](https://github.com/flamingo-stack/openframe-oss-lib/blob/main/src/main/java/com/openframe/stream/deserializer/Microsoft365AuditEventDeserializer.java) — template this deserializer mirrors -- [`EventTypeMapper`](https://github.com/flamingo-stack/openframe-oss-lib/blob/main/src/main/java/com/openframe/stream/mapping/EventTypeMapper.java) — `eventType` → `UnifiedEventType` registry -- [`SourceEventTypes.GoogleWorkspace`](https://github.com/flamingo-stack/openframe-oss-lib/blob/main/src/main/java/com/openframe/stream/mapping/SourceEventTypes.java) — `eventType` string constants -- [`MessageType`](https://github.com/flamingo-stack/openframe-oss-lib/blob/main/src/main/java/com/openframe/data/model/enums/MessageType.java) — enum of all routable message types diff --git a/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/.Microsoft365AuditEventDeserializer.md b/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/.Microsoft365AuditEventDeserializer.md deleted file mode 100644 index 080c5fafe..000000000 --- a/openframe-stream-service-core/src/main/java/com/openframe/stream/deserializer/.Microsoft365AuditEventDeserializer.md +++ /dev/null @@ -1,53 +0,0 @@ - -A Spring component that deserializes Microsoft 365 Entra directory audit events polled from Graph `auditLogs/directoryAudits`. Unlike the CDC-sourced deserializers in this package, the events are hand-built by the poller (not Debezium change capture), arrive pre-enriched with tenant/organization fields, and carry no agent reference — routed via `DataEnrichmentServiceType.PRE_ENRICHED`. - -## Key Components - -- **`getType()`** — Returns `MessageType.MICROSOFT_365_AUDIT_EVENT` -- **`deserialize()`** — Maps the flat `after` payload built by the poller into a `DeserializedDebeziumMessage`; returns `null` when `after` is null -- **`resolveEventType()`** — Failure-first classification: a `result` of `failure`/`timeout` (case-insensitive) always maps to `UnifiedEventType.M365_AUDIT_FAILURE` regardless of `category`; otherwise delegates to `EventTypeMapper` and falls back to `M365_AUDIT_OTHER` for unmapped categories (never `UNKNOWN`) -- **`getEventTimestamp()`** — Parses `activityDateTime` (RFC3339/ISO-8601); falls back to the Debezium payload's processing `timestamp` when absent or unparseable -- **`buildDetails()`** — Verbatim `initiatedBy`, `targetResources`, `additionalDetails` sub-objects, plus flat `connectionId`/`connectionName` (multi-connection orgs) - -## Field Map - -| `DeserializedDebeziumMessage` field | Source in `after` | -|---|---| -| `toolEventId` | `auditId` (Graph audit record id — idempotent across poller replay) | -| `eventTimestamp` | `activityDateTime`, fallback processing timestamp | -| `message` | `activityDisplayName` | -| `sourceEventType` | `category` | -| `userId` | `initiatedBy.user.userPrincipalName` | -| `tenantId` / `organizationId` / `organizationName` | passthrough | -| `agentId` | always `null` | -| `isVisible` / `skipProcessing` | always `true` / `false` | - -## Usage Example - -```java -ObjectMapper mapper = new ObjectMapper(); -Microsoft365AuditEventDeserializer deserializer = new Microsoft365AuditEventDeserializer(mapper); - -// Spring wires this as the strategy registered for MessageType.MICROSOFT_365_AUDIT_EVENT -MessageType type = deserializer.getType(); -// → MessageType.MICROSOFT_365_AUDIT_EVENT - -JsonNode after = mapper.readTree(""" - { - "auditId": "Directory_abc_123", "activityDateTime": "2026-07-30T10:00:00Z", - "activityDisplayName": "Add user", "category": "UserManagement", "result": "success", - "initiatedBy": {"user": {"userPrincipalName": "admin@x.com"}}, - "tenantId": "tenant-1", "organizationId": "org-uuid-1", "organizationName": "Acme Org" - } - """); - -DeserializedDebeziumMessage result = deserializer.deserialize(debeziumMessageWith(after), type); -// → unifiedEventType = M365_USER_MANAGEMENT, userId = "admin@x.com", toolEventId = "Directory_abc_123" -``` - -## Related Files - -- [`GoogleWorkspaceAuditEventDeserializer`](https://github.com/flamingo-stack/openframe-oss-lib/blob/main/src/main/java/com/openframe/stream/deserializer/GoogleWorkspaceAuditEventDeserializer.java) — sibling deserializer for Google Workspace directory audit events, same poller-fed / pre-enriched / 1:1 pattern -- [`EventTypeMapper`](https://github.com/flamingo-stack/openframe-oss-lib/blob/main/src/main/java/com/openframe/stream/mapping/EventTypeMapper.java) — `category` → `UnifiedEventType` registry -- [`SourceEventTypes.Microsoft365`](https://github.com/flamingo-stack/openframe-oss-lib/blob/main/src/main/java/com/openframe/stream/mapping/SourceEventTypes.java) — `category` string constants -- [`MessageType`](https://github.com/flamingo-stack/openframe-oss-lib/blob/main/src/main/java/com/openframe/data/model/enums/MessageType.java) — enum of all routable message types From bd12939d52835a45a96d79655a9b5963939c7a93 Mon Sep 17 00:00:00 2001 From: andrii Date: Wed, 12 Aug 2026 11:45:17 +0200 Subject: [PATCH 9/9] docs(stream): clarify multiMessageValue handling in GWS deserializer The Javadoc listed multiMessageValue among the Reports API parameter value keys but left it out of the "value ?? boolValue ?? intValue ?? multiValue" guidance, so a consumer following that chain could silently drop it. Spell out that multiMessageValue is intentionally excluded from the scalar fallback and needs recursive rendering of its nested parameter[] objects. Co-Authored-By: Claude Opus 5 --- .../GoogleWorkspaceAuditEventDeserializer.java | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) 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 index 9afb5f9db..a46c3565a 100644 --- 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 @@ -38,10 +38,14 @@ *

* {@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} 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}; nothing may assume {@code value} alone. + * {@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