From 397d4ad651188f336e88252f5c0604c5101d6dba Mon Sep 17 00:00:00 2001 From: Hryhorii Chuhuievets Date: Mon, 7 Sep 2026 13:25:49 +0300 Subject: [PATCH 1/5] Use nickname in log event instead of hostname --- .../data/model/redis/CachedMachineInfo.java | 1 + .../redis/MachineIdCacheService.java | 5 ++- .../IntegratedToolDataEnrichmentService.java | 8 +++- .../service/MachineDisplayNameResolver.java | 18 ++++++++ .../service/rmm/RmmEnrichmentService.java | 9 +++- .../MachineDisplayNameResolverTest.java | 41 +++++++++++++++++++ .../service/RmmEnrichmentServiceTest.java | 26 +++++++++--- ...riptExecutedEnrichmentIntegrationTest.java | 8 ++-- 8 files changed, 102 insertions(+), 14 deletions(-) create mode 100644 openframe-stream-service-core/src/main/java/com/openframe/stream/service/MachineDisplayNameResolver.java create mode 100644 openframe-stream-service-core/src/test/java/com/openframe/stream/service/MachineDisplayNameResolverTest.java diff --git a/openframe-data-mongo-common/src/main/java/com/openframe/data/model/redis/CachedMachineInfo.java b/openframe-data-mongo-common/src/main/java/com/openframe/data/model/redis/CachedMachineInfo.java index 99b462b0e9..2ba498756f 100644 --- a/openframe-data-mongo-common/src/main/java/com/openframe/data/model/redis/CachedMachineInfo.java +++ b/openframe-data-mongo-common/src/main/java/com/openframe/data/model/redis/CachedMachineInfo.java @@ -18,6 +18,7 @@ public class CachedMachineInfo implements Serializable { private String machineId; private String hostname; + private String nickname; private String organizationId; } diff --git a/openframe-data-mongo-sync/src/main/java/com/openframe/data/repository/redis/MachineIdCacheService.java b/openframe-data-mongo-sync/src/main/java/com/openframe/data/repository/redis/MachineIdCacheService.java index c4403ce803..1acd710c9a 100644 --- a/openframe-data-mongo-sync/src/main/java/com/openframe/data/repository/redis/MachineIdCacheService.java +++ b/openframe-data-mongo-sync/src/main/java/com/openframe/data/repository/redis/MachineIdCacheService.java @@ -30,7 +30,7 @@ public class MachineIdCacheService { /** * Get cached machine info from cache or database by agent ID - * Returns only essential fields (machineId, hostname, organizationId) +` * Returns only essential fields (machineId, hostname, nickname, organizationId) * * @param agentId the agent ID * @return the CachedMachineInfo object, or null if not found @@ -46,6 +46,7 @@ public CachedMachineInfo getMachine(String agentId) { .map(machine -> new CachedMachineInfo( machine.getMachineId(), machine.getHostname(), + machine.getNickname(), machine.getOrganizationId() )) .orElse(null); @@ -66,6 +67,7 @@ public CachedMachineInfo getMachine(String tenantId, ToolType toolType, String a .map(machine -> new CachedMachineInfo( machine.getMachineId(), machine.getHostname(), + machine.getNickname(), machine.getOrganizationId() )) .orElse(null); @@ -90,6 +92,7 @@ public CachedMachineInfo getMachineByMachineId(String machineId) { .map(machine -> new CachedMachineInfo( machine.getMachineId(), machine.getHostname(), + machine.getNickname(), machine.getOrganizationId() )) .orElse(null); diff --git a/openframe-stream-service-core/src/main/java/com/openframe/stream/service/IntegratedToolDataEnrichmentService.java b/openframe-stream-service-core/src/main/java/com/openframe/stream/service/IntegratedToolDataEnrichmentService.java index 21831b593a..7cdb034ba0 100644 --- a/openframe-stream-service-core/src/main/java/com/openframe/stream/service/IntegratedToolDataEnrichmentService.java +++ b/openframe-stream-service-core/src/main/java/com/openframe/stream/service/IntegratedToolDataEnrichmentService.java @@ -22,13 +22,16 @@ public class IntegratedToolDataEnrichmentService implements DataEnrichmentServic private final MachineIdCacheService machineIdCacheService; private final ClusterTenantIdResolver clusterTenantIdResolver; private final TenantIdProvider tenantIdProvider; + private final MachineDisplayNameResolver machineDisplayNameResolver; public IntegratedToolDataEnrichmentService(MachineIdCacheService machineIdCacheService, @Autowired(required = false) ClusterTenantIdResolver clusterTenantIdResolver, - TenantIdProvider tenantIdProvider) { + TenantIdProvider tenantIdProvider, + MachineDisplayNameResolver machineDisplayNameResolver) { this.machineIdCacheService = machineIdCacheService; this.clusterTenantIdResolver = clusterTenantIdResolver; this.tenantIdProvider = tenantIdProvider; + this.machineDisplayNameResolver = machineDisplayNameResolver; } @Override @@ -53,8 +56,9 @@ private void enrichFromMachine(DeserializedDebeziumMessage message, IntegratedTo log.warn("Machine ID not found for agent: {}", agentId); return; } + String displayName = machineDisplayNameResolver.resolveDisplayName(machine); enriched.setMachineId(machine.getMachineId()); - enriched.setHostname(machine.getHostname()); + enriched.setHostname(displayName); CachedOrganizationInfo organization = machineIdCacheService.getOrganization(machine.getOrganizationId()); log.debug("Found machine ID {} for agent {} (organization {})", diff --git a/openframe-stream-service-core/src/main/java/com/openframe/stream/service/MachineDisplayNameResolver.java b/openframe-stream-service-core/src/main/java/com/openframe/stream/service/MachineDisplayNameResolver.java new file mode 100644 index 0000000000..a773e1ead7 --- /dev/null +++ b/openframe-stream-service-core/src/main/java/com/openframe/stream/service/MachineDisplayNameResolver.java @@ -0,0 +1,18 @@ +package com.openframe.stream.service; + +import com.openframe.data.model.redis.CachedMachineInfo; +import org.springframework.stereotype.Component; + +import static org.springframework.util.StringUtils.hasText; + +@Component +public class MachineDisplayNameResolver { + + public String resolveDisplayName(CachedMachineInfo machine) { + String nickname = machine.getNickname(); + if (hasText(nickname)) { + return nickname; + } + return machine.getHostname(); + } +} diff --git a/openframe-stream-service-core/src/main/java/com/openframe/stream/service/rmm/RmmEnrichmentService.java b/openframe-stream-service-core/src/main/java/com/openframe/stream/service/rmm/RmmEnrichmentService.java index 32bdc3b5c5..e35f7cea30 100644 --- a/openframe-stream-service-core/src/main/java/com/openframe/stream/service/rmm/RmmEnrichmentService.java +++ b/openframe-stream-service-core/src/main/java/com/openframe/stream/service/rmm/RmmEnrichmentService.java @@ -10,6 +10,7 @@ import com.openframe.stream.service.ClusterTenantIdResolver; import com.openframe.stream.service.DataEnrichmentService; import com.openframe.stream.service.IntegratedToolDataEnrichmentService; +import com.openframe.stream.service.MachineDisplayNameResolver; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; @@ -34,13 +35,16 @@ public class RmmEnrichmentService implements DataEnrichmentService Date: Tue, 8 Sep 2026 17:56:05 +0300 Subject: [PATCH 2/5] Add device Nickname to LogEvent --- .../openframe/api/dto/audit/LogDetails.java | 1 + .../com/openframe/api/dto/audit/LogEvent.java | 1 + .../com/openframe/api/service/LogService.java | 2 + .../src/main/resources/schema/log.graphqls | 2 + .../data/cassandra/model/UnifiedLogEvent.java | 3 ++ .../kafka/model/IntegratedToolEvent.java | 1 + .../data/pinot/model/LogProjection.java | 1 + .../repository/PinotClientLogRepository.java | 5 ++- .../PinotClientLogRepositoryTest.java | 26 +++++++----- .../dto/audit/LogDetailsResponse.java | 3 ++ .../external/dto/audit/LogResponse.java | 3 ++ .../openframe/external/mapper/LogMapper.java | 2 + .../resources/pinot/config/schema-logs.json | 1 + .../config/table-config-logs-realtime.json | 12 ++++++ .../DebeziumCassandraMessageHandler.java | 1 + .../handler/DebeziumKafkaMessageHandler.java | 1 + .../debezium/IntegratedToolEnrichedData.java | 1 + .../IntegratedToolDataEnrichmentService.java | 9 ++-- .../service/MachineDisplayNameResolver.java | 18 -------- .../service/rmm/RmmEnrichmentService.java | 10 ++--- .../MachineDisplayNameResolverTest.java | 41 ------------------- .../service/RmmEnrichmentServiceTest.java | 24 ++++++++--- ...riptExecutedEnrichmentIntegrationTest.java | 4 +- 23 files changed, 80 insertions(+), 92 deletions(-) delete mode 100644 openframe-stream-service-core/src/main/java/com/openframe/stream/service/MachineDisplayNameResolver.java delete mode 100644 openframe-stream-service-core/src/test/java/com/openframe/stream/service/MachineDisplayNameResolverTest.java diff --git a/openframe-api-lib/src/main/java/com/openframe/api/dto/audit/LogDetails.java b/openframe-api-lib/src/main/java/com/openframe/api/dto/audit/LogDetails.java index f1a64b3e31..9fa5b4b801 100644 --- a/openframe-api-lib/src/main/java/com/openframe/api/dto/audit/LogDetails.java +++ b/openframe-api-lib/src/main/java/com/openframe/api/dto/audit/LogDetails.java @@ -21,6 +21,7 @@ public class LogDetails { private String userId; private String deviceId; private String hostname; + private String nickname; private String organizationId; private String organizationName; private String summary; diff --git a/openframe-api-lib/src/main/java/com/openframe/api/dto/audit/LogEvent.java b/openframe-api-lib/src/main/java/com/openframe/api/dto/audit/LogEvent.java index 2310a06956..eb6819ac33 100644 --- a/openframe-api-lib/src/main/java/com/openframe/api/dto/audit/LogEvent.java +++ b/openframe-api-lib/src/main/java/com/openframe/api/dto/audit/LogEvent.java @@ -21,6 +21,7 @@ public class LogEvent { private String userId; private String deviceId; private String hostname; + private String nickname; private String organizationId; private String organizationName; private String summary; diff --git a/openframe-api-lib/src/main/java/com/openframe/api/service/LogService.java b/openframe-api-lib/src/main/java/com/openframe/api/service/LogService.java index 6c1137f987..e5ebcff526 100644 --- a/openframe-api-lib/src/main/java/com/openframe/api/service/LogService.java +++ b/openframe-api-lib/src/main/java/com/openframe/api/service/LogService.java @@ -203,6 +203,7 @@ private LogEvent mapToLogEvent(LogProjection log) { .userId(log.userId) .deviceId(log.deviceId) .hostname(log.hostname) + .nickname(log.nickname) .organizationId(log.organizationId) .organizationName(log.organizationName) .build(); @@ -222,6 +223,7 @@ private LogDetails mapToLogDetails(UnifiedLogEvent logEvent) { .userId(logEvent.getUserId()) .deviceId(logEvent.getDeviceId()) .hostname(logEvent.getHostname()) + .nickname(logEvent.getNickname()) .organizationId(logEvent.getOrganizationId()) .organizationName(logEvent.getOrganizationName()) .summary(logEvent.getMessage()) diff --git a/openframe-api-service-core/src/main/resources/schema/log.graphqls b/openframe-api-service-core/src/main/resources/schema/log.graphqls index 804b0d0efa..1c63a9b982 100644 --- a/openframe-api-service-core/src/main/resources/schema/log.graphqls +++ b/openframe-api-service-core/src/main/resources/schema/log.graphqls @@ -78,6 +78,7 @@ type LogEvent implements Node { userId: String deviceId: String hostname: String + nickname: String organizationId: String organizationName: String summary: String @@ -95,6 +96,7 @@ type LogDetails implements Node { userId: String deviceId: String hostname: String + nickname: String organizationId: String organizationName: String message: String diff --git a/openframe-data-cassandra/src/main/java/com/openframe/data/cassandra/model/UnifiedLogEvent.java b/openframe-data-cassandra/src/main/java/com/openframe/data/cassandra/model/UnifiedLogEvent.java index 8cedd14df9..8b06f29dce 100644 --- a/openframe-data-cassandra/src/main/java/com/openframe/data/cassandra/model/UnifiedLogEvent.java +++ b/openframe-data-cassandra/src/main/java/com/openframe/data/cassandra/model/UnifiedLogEvent.java @@ -35,6 +35,9 @@ public class UnifiedLogEvent { @Column("hostname") private String hostname; + @Column("nickname") + private String nickname; + /** * Organization ID associated with the event. */ diff --git a/openframe-data-kafka/src/main/java/com/openframe/kafka/model/IntegratedToolEvent.java b/openframe-data-kafka/src/main/java/com/openframe/kafka/model/IntegratedToolEvent.java index 7518229763..52daf922a3 100644 --- a/openframe-data-kafka/src/main/java/com/openframe/kafka/model/IntegratedToolEvent.java +++ b/openframe-data-kafka/src/main/java/com/openframe/kafka/model/IntegratedToolEvent.java @@ -9,6 +9,7 @@ public class IntegratedToolEvent implements KafkaMessage { private String userId; private String deviceId; private String hostname; + private String nickname; private String organizationId; private String organizationName; private String ingestDay; diff --git a/openframe-data-pinot/src/main/java/com/openframe/data/pinot/model/LogProjection.java b/openframe-data-pinot/src/main/java/com/openframe/data/pinot/model/LogProjection.java index 9fcad2f08d..efee03883d 100644 --- a/openframe-data-pinot/src/main/java/com/openframe/data/pinot/model/LogProjection.java +++ b/openframe-data-pinot/src/main/java/com/openframe/data/pinot/model/LogProjection.java @@ -21,6 +21,7 @@ public class LogProjection { public String userId; public String deviceId; public String hostname; + public String nickname; public String organizationId; public String organizationName; public String summary; diff --git a/openframe-data-pinot/src/main/java/com/openframe/data/pinot/repository/PinotClientLogRepository.java b/openframe-data-pinot/src/main/java/com/openframe/data/pinot/repository/PinotClientLogRepository.java index ad7faa3c2e..37304ed6cc 100644 --- a/openframe-data-pinot/src/main/java/com/openframe/data/pinot/repository/PinotClientLogRepository.java +++ b/openframe-data-pinot/src/main/java/com/openframe/data/pinot/repository/PinotClientLogRepository.java @@ -45,7 +45,7 @@ public List findLogs(String tenantId, LocalDate startDate, LocalD List severities, List organizationIds, String deviceId, String cursor, int limit, String sortField, String sortDirection) { PinotQueryBuilder queryBuilder = new PinotQueryBuilder(logsTable, tenantId) - .select("toolEventId", "ingestDay", "toolType", "eventType", "severity", "userId", "deviceId", "hostname", "organizationId", "organizationName", "summary", "eventTimestamp") + .select("toolEventId", "ingestDay", "toolType", "eventType", "severity", "userId", "deviceId", "hostname", "nickname", "organizationId", "organizationName", "summary", "eventTimestamp") .whereDateRange("eventTimestamp", startDate, endDate) .whereTimestampRange("eventTimestamp", timestampFrom, timestampTo) .whereIn("toolType", toolTypes) @@ -66,7 +66,7 @@ public List searchLogs(String tenantId, LocalDate startDate, Loca List severities, List organizationIds, String deviceId, String searchTerm, String cursor, int limit, String sortField, String sortDirection) { PinotQueryBuilder queryBuilder = new PinotQueryBuilder(logsTable, tenantId) - .select("toolEventId", "ingestDay", "toolType", "eventType", "severity", "userId", "deviceId", "hostname", "organizationId", "organizationName", "summary", "eventTimestamp") + .select("toolEventId", "ingestDay", "toolType", "eventType", "severity", "userId", "deviceId", "hostname", "nickname", "organizationId", "organizationName", "summary", "eventTimestamp") .whereDateRange("eventTimestamp", startDate, endDate) .whereTimestampRange("eventTimestamp", timestampFrom, timestampTo) .whereIn("toolType", toolTypes) @@ -191,6 +191,7 @@ private List executeLogQuery(String query) { projection.userId = resultSet.getString(rowIndex, columnIndexMap.get("userId")); projection.deviceId = resultSet.getString(rowIndex, columnIndexMap.get("deviceId")); projection.hostname = resultSet.getString(rowIndex, columnIndexMap.get("hostname")); + projection.nickname = resultSet.getString(rowIndex, columnIndexMap.get("nickname")); projection.organizationId = resultSet.getString(rowIndex, columnIndexMap.get("organizationId")); projection.organizationName = resultSet.getString(rowIndex, columnIndexMap.get("organizationName")); projection.summary = resultSet.getString(rowIndex, columnIndexMap.get("summary")); diff --git a/openframe-data-pinot/src/test/java/com/openframe/data/pinot/repository/PinotClientLogRepositoryTest.java b/openframe-data-pinot/src/test/java/com/openframe/data/pinot/repository/PinotClientLogRepositoryTest.java index d0536f3d03..adc77c609f 100644 --- a/openframe-data-pinot/src/test/java/com/openframe/data/pinot/repository/PinotClientLogRepositoryTest.java +++ b/openframe-data-pinot/src/test/java/com/openframe/data/pinot/repository/PinotClientLogRepositoryTest.java @@ -175,12 +175,12 @@ class ResultMapping { @Test @DisplayName("findLogs maps ResultSet rows to LogProjection objects") void mapsLogProjection() { - // 12 columns as declared in findLogs: toolEventId, ingestDay, toolType, eventType, - // severity, userId, deviceId, hostname, organizationId, organizationName, summary, eventTimestamp + // 13 columns as declared in findLogs: toolEventId, ingestDay, toolType, eventType, severity, + // userId, deviceId, hostname, nickname, organizationId, organizationName, summary, eventTimestamp when(pinotConnection.execute(anyString())).thenReturn(resultSetGroup); when(resultSetGroup.getResultSet(0)).thenReturn(resultSet); when(resultSet.getRowCount()).thenReturn(1); - when(resultSet.getColumnCount()).thenReturn(12); + when(resultSet.getColumnCount()).thenReturn(13); when(resultSet.getColumnName(0)).thenReturn("toolEventId"); when(resultSet.getColumnName(1)).thenReturn("ingestDay"); when(resultSet.getColumnName(2)).thenReturn("toolType"); @@ -189,10 +189,11 @@ void mapsLogProjection() { when(resultSet.getColumnName(5)).thenReturn("userId"); when(resultSet.getColumnName(6)).thenReturn("deviceId"); when(resultSet.getColumnName(7)).thenReturn("hostname"); - when(resultSet.getColumnName(8)).thenReturn("organizationId"); - when(resultSet.getColumnName(9)).thenReturn("organizationName"); - when(resultSet.getColumnName(10)).thenReturn("summary"); - when(resultSet.getColumnName(11)).thenReturn("eventTimestamp"); + when(resultSet.getColumnName(8)).thenReturn("nickname"); + when(resultSet.getColumnName(9)).thenReturn("organizationId"); + when(resultSet.getColumnName(10)).thenReturn("organizationName"); + when(resultSet.getColumnName(11)).thenReturn("summary"); + when(resultSet.getColumnName(12)).thenReturn("eventTimestamp"); when(resultSet.getString(0, 0)).thenReturn("evt-1"); when(resultSet.getString(0, 1)).thenReturn("2026-01-15"); @@ -202,10 +203,11 @@ void mapsLogProjection() { when(resultSet.getString(0, 5)).thenReturn("user-1"); when(resultSet.getString(0, 6)).thenReturn("dev-1"); when(resultSet.getString(0, 7)).thenReturn("host-1"); - when(resultSet.getString(0, 8)).thenReturn("org-1"); - when(resultSet.getString(0, 9)).thenReturn("Org One"); - when(resultSet.getString(0, 10)).thenReturn("Device came online"); - when(resultSet.getLong(0, 11)).thenReturn(1700000000000L); + when(resultSet.getString(0, 8)).thenReturn("Reception iMac"); + when(resultSet.getString(0, 9)).thenReturn("org-1"); + when(resultSet.getString(0, 10)).thenReturn("Org One"); + when(resultSet.getString(0, 11)).thenReturn("Device came online"); + when(resultSet.getLong(0, 12)).thenReturn(1700000000000L); List result = repository.findLogs(TENANT_ID, START, END, null, null, List.of(), List.of(), List.of(), List.of(), null, null, 100, "eventTimestamp", "DESC"); @@ -215,6 +217,8 @@ void mapsLogProjection() { assertEquals("evt-1", p.toolEventId); assertEquals("FLEET_MDM", p.toolType); assertEquals("INFO", p.severity); + assertEquals("host-1", p.hostname); + assertEquals("Reception iMac", p.nickname); assertEquals("Org One", p.organizationName); assertNotNull(p.eventTimestamp); } diff --git a/openframe-external-api-service-core/src/main/java/com/openframe/external/dto/audit/LogDetailsResponse.java b/openframe-external-api-service-core/src/main/java/com/openframe/external/dto/audit/LogDetailsResponse.java index 4276852b35..5830d49e78 100644 --- a/openframe-external-api-service-core/src/main/java/com/openframe/external/dto/audit/LogDetailsResponse.java +++ b/openframe-external-api-service-core/src/main/java/com/openframe/external/dto/audit/LogDetailsResponse.java @@ -40,6 +40,9 @@ public class LogDetailsResponse { @Schema(description = "Hostname of the device associated with the event") private String hostname; + @Schema(description = "Nickname of the device associated with the event") + private String nickname; + @Schema(description = "Customer id associated with the event") private String customerId; diff --git a/openframe-external-api-service-core/src/main/java/com/openframe/external/dto/audit/LogResponse.java b/openframe-external-api-service-core/src/main/java/com/openframe/external/dto/audit/LogResponse.java index 4e85abd3d3..8ed9cb55ba 100644 --- a/openframe-external-api-service-core/src/main/java/com/openframe/external/dto/audit/LogResponse.java +++ b/openframe-external-api-service-core/src/main/java/com/openframe/external/dto/audit/LogResponse.java @@ -39,6 +39,9 @@ public class LogResponse { @Schema(description = "Hostname of the device associated with the event") private String hostname; + @Schema(description = "Nickname of the device associated with the event") + private String nickname; + @Schema(description = "Customer id associated with the event") private String customerId; diff --git a/openframe-external-api-service-core/src/main/java/com/openframe/external/mapper/LogMapper.java b/openframe-external-api-service-core/src/main/java/com/openframe/external/mapper/LogMapper.java index 48504c0d68..b7854c2f99 100644 --- a/openframe-external-api-service-core/src/main/java/com/openframe/external/mapper/LogMapper.java +++ b/openframe-external-api-service-core/src/main/java/com/openframe/external/mapper/LogMapper.java @@ -32,6 +32,7 @@ public LogResponse toLogResponse(LogEvent logEvent) { .userId(logEvent.getUserId()) .deviceId(logEvent.getDeviceId()) .hostname(logEvent.getHostname()) + .nickname(logEvent.getNickname()) .customerId(logEvent.getOrganizationId()) .customerName(logEvent.getOrganizationName()) .summary(logEvent.getSummary()) @@ -90,6 +91,7 @@ public LogDetailsResponse toLogDetailsResponse(LogDetails logDetails) { .userId(logDetails.getUserId()) .deviceId(logDetails.getDeviceId()) .hostname(logDetails.getHostname()) + .nickname(logDetails.getNickname()) .customerId(logDetails.getOrganizationId()) .customerName(logDetails.getOrganizationName()) .summary(logDetails.getSummary()) diff --git a/openframe-pinot-initializer/src/main/resources/pinot/config/schema-logs.json b/openframe-pinot-initializer/src/main/resources/pinot/config/schema-logs.json index 69e102445a..4d4eb1ac92 100644 --- a/openframe-pinot-initializer/src/main/resources/pinot/config/schema-logs.json +++ b/openframe-pinot-initializer/src/main/resources/pinot/config/schema-logs.json @@ -7,6 +7,7 @@ { "name": "userId", "dataType": "STRING", "singleValueField": true }, { "name": "deviceId", "dataType": "STRING", "singleValueField": true }, { "name": "hostname", "dataType": "STRING", "singleValueField": true }, + { "name": "nickname", "dataType": "STRING", "singleValueField": true }, { "name": "organizationId", "dataType": "STRING", "singleValueField": true }, { "name": "organizationName", "dataType": "STRING", "singleValueField": true }, { "name": "ingestDay", "dataType": "STRING", "singleValueField": true }, diff --git a/openframe-pinot-initializer/src/main/resources/pinot/config/table-config-logs-realtime.json b/openframe-pinot-initializer/src/main/resources/pinot/config/table-config-logs-realtime.json index 5a56168920..e08fc35a8a 100644 --- a/openframe-pinot-initializer/src/main/resources/pinot/config/table-config-logs-realtime.json +++ b/openframe-pinot-initializer/src/main/resources/pinot/config/table-config-logs-realtime.json @@ -35,6 +35,7 @@ "summary", "userId", "hostname", + "nickname", "organizationName" ], "invertedIndexColumns": [ @@ -106,6 +107,17 @@ "enableQueryCache": "true", "caseSensitive": "false" } + }, + { + "name": "nickname", + "encodingType": "RAW", + "indexTypes": [ + "TEXT" + ], + "properties": { + "enableQueryCache": "true", + "caseSensitive": "false" + } } ], "routing": { diff --git a/openframe-stream-service-core/src/main/java/com/openframe/stream/handler/DebeziumCassandraMessageHandler.java b/openframe-stream-service-core/src/main/java/com/openframe/stream/handler/DebeziumCassandraMessageHandler.java index 49d97bd2c8..19a11aa17b 100644 --- a/openframe-stream-service-core/src/main/java/com/openframe/stream/handler/DebeziumCassandraMessageHandler.java +++ b/openframe-stream-service-core/src/main/java/com/openframe/stream/handler/DebeziumCassandraMessageHandler.java @@ -61,6 +61,7 @@ protected UnifiedLogEvent transform(DeserializedDebeziumMessage debeziumMessage, logEvent.setUserId(enrichedData.getUserId()); logEvent.setDeviceId(enrichedData.getMachineId()); logEvent.setHostname(enrichedData.getHostname()); + logEvent.setNickname(enrichedData.getNickname()); logEvent.setOrganizationId(enrichedData.getOrganizationId()); logEvent.setOrganizationName(enrichedData.getOrganizationName()); logEvent.setSeverity(debeziumMessage.getUnifiedEventType().getSeverity().name()); diff --git a/openframe-stream-service-core/src/main/java/com/openframe/stream/handler/DebeziumKafkaMessageHandler.java b/openframe-stream-service-core/src/main/java/com/openframe/stream/handler/DebeziumKafkaMessageHandler.java index 9b1c1666b8..15f1ccb4ae 100644 --- a/openframe-stream-service-core/src/main/java/com/openframe/stream/handler/DebeziumKafkaMessageHandler.java +++ b/openframe-stream-service-core/src/main/java/com/openframe/stream/handler/DebeziumKafkaMessageHandler.java @@ -35,6 +35,7 @@ protected IntegratedToolEvent transform(DeserializedDebeziumMessage debeziumMess message.setUserId(enrichedData.getUserId()); message.setDeviceId(enrichedData.getMachineId()); message.setHostname(enrichedData.getHostname()); + message.setNickname(enrichedData.getNickname()); message.setOrganizationId(enrichedData.getOrganizationId()); message.setOrganizationName(enrichedData.getOrganizationName()); message.setIngestDay(debeziumMessage.getIngestDay()); diff --git a/openframe-stream-service-core/src/main/java/com/openframe/stream/model/fleet/debezium/IntegratedToolEnrichedData.java b/openframe-stream-service-core/src/main/java/com/openframe/stream/model/fleet/debezium/IntegratedToolEnrichedData.java index d9ba7d2e46..5a897bc06f 100644 --- a/openframe-stream-service-core/src/main/java/com/openframe/stream/model/fleet/debezium/IntegratedToolEnrichedData.java +++ b/openframe-stream-service-core/src/main/java/com/openframe/stream/model/fleet/debezium/IntegratedToolEnrichedData.java @@ -7,6 +7,7 @@ public class IntegratedToolEnrichedData { private String machineId; private String hostname; + private String nickname; private String organizationId; private String organizationName; private String userId; diff --git a/openframe-stream-service-core/src/main/java/com/openframe/stream/service/IntegratedToolDataEnrichmentService.java b/openframe-stream-service-core/src/main/java/com/openframe/stream/service/IntegratedToolDataEnrichmentService.java index 7cdb034ba0..f11cb3919b 100644 --- a/openframe-stream-service-core/src/main/java/com/openframe/stream/service/IntegratedToolDataEnrichmentService.java +++ b/openframe-stream-service-core/src/main/java/com/openframe/stream/service/IntegratedToolDataEnrichmentService.java @@ -22,16 +22,13 @@ public class IntegratedToolDataEnrichmentService implements DataEnrichmentServic private final MachineIdCacheService machineIdCacheService; private final ClusterTenantIdResolver clusterTenantIdResolver; private final TenantIdProvider tenantIdProvider; - private final MachineDisplayNameResolver machineDisplayNameResolver; public IntegratedToolDataEnrichmentService(MachineIdCacheService machineIdCacheService, @Autowired(required = false) ClusterTenantIdResolver clusterTenantIdResolver, - TenantIdProvider tenantIdProvider, - MachineDisplayNameResolver machineDisplayNameResolver) { + TenantIdProvider tenantIdProvider) { this.machineIdCacheService = machineIdCacheService; this.clusterTenantIdResolver = clusterTenantIdResolver; this.tenantIdProvider = tenantIdProvider; - this.machineDisplayNameResolver = machineDisplayNameResolver; } @Override @@ -56,9 +53,9 @@ private void enrichFromMachine(DeserializedDebeziumMessage message, IntegratedTo log.warn("Machine ID not found for agent: {}", agentId); return; } - String displayName = machineDisplayNameResolver.resolveDisplayName(machine); enriched.setMachineId(machine.getMachineId()); - enriched.setHostname(displayName); + enriched.setHostname(machine.getHostname()); + enriched.setNickname(machine.getNickname()); CachedOrganizationInfo organization = machineIdCacheService.getOrganization(machine.getOrganizationId()); log.debug("Found machine ID {} for agent {} (organization {})", diff --git a/openframe-stream-service-core/src/main/java/com/openframe/stream/service/MachineDisplayNameResolver.java b/openframe-stream-service-core/src/main/java/com/openframe/stream/service/MachineDisplayNameResolver.java deleted file mode 100644 index a773e1ead7..0000000000 --- a/openframe-stream-service-core/src/main/java/com/openframe/stream/service/MachineDisplayNameResolver.java +++ /dev/null @@ -1,18 +0,0 @@ -package com.openframe.stream.service; - -import com.openframe.data.model.redis.CachedMachineInfo; -import org.springframework.stereotype.Component; - -import static org.springframework.util.StringUtils.hasText; - -@Component -public class MachineDisplayNameResolver { - - public String resolveDisplayName(CachedMachineInfo machine) { - String nickname = machine.getNickname(); - if (hasText(nickname)) { - return nickname; - } - return machine.getHostname(); - } -} diff --git a/openframe-stream-service-core/src/main/java/com/openframe/stream/service/rmm/RmmEnrichmentService.java b/openframe-stream-service-core/src/main/java/com/openframe/stream/service/rmm/RmmEnrichmentService.java index e35f7cea30..cf792f2abf 100644 --- a/openframe-stream-service-core/src/main/java/com/openframe/stream/service/rmm/RmmEnrichmentService.java +++ b/openframe-stream-service-core/src/main/java/com/openframe/stream/service/rmm/RmmEnrichmentService.java @@ -10,7 +10,6 @@ import com.openframe.stream.service.ClusterTenantIdResolver; import com.openframe.stream.service.DataEnrichmentService; import com.openframe.stream.service.IntegratedToolDataEnrichmentService; -import com.openframe.stream.service.MachineDisplayNameResolver; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; @@ -35,16 +34,13 @@ public class RmmEnrichmentService implements DataEnrichmentService Date: Tue, 8 Sep 2026 19:54:25 +0300 Subject: [PATCH 3/5] Changes after initial code review --- .../repository/AbstractPinotRepository.java | 11 ++ .../repository/PinotClientLogRepository.java | 37 ++++-- .../pinot/repository/PinotQueryBuilder.java | 2 + .../PinotClientLogRepositoryTest.java | 86 ++++++++++++++ .../config/pinot/PinotConfigInitializer.java | 46 +++++++- .../resources/pinot/config/schema-logs.json | 2 +- .../LogEventNicknamePropagationTest.java | 105 +++++++++++++++++ ...tegratedToolDataEnrichmentServiceTest.java | 109 ++++++++++++++++++ ...riptExecutedEnrichmentIntegrationTest.java | 7 +- 9 files changed, 387 insertions(+), 18 deletions(-) create mode 100644 openframe-stream-service-core/src/test/java/com/openframe/stream/handler/LogEventNicknamePropagationTest.java create mode 100644 openframe-stream-service-core/src/test/java/com/openframe/stream/service/IntegratedToolDataEnrichmentServiceTest.java diff --git a/openframe-data-pinot/src/main/java/com/openframe/data/pinot/repository/AbstractPinotRepository.java b/openframe-data-pinot/src/main/java/com/openframe/data/pinot/repository/AbstractPinotRepository.java index db9257f0b8..04967f0c79 100644 --- a/openframe-data-pinot/src/main/java/com/openframe/data/pinot/repository/AbstractPinotRepository.java +++ b/openframe-data-pinot/src/main/java/com/openframe/data/pinot/repository/AbstractPinotRepository.java @@ -112,6 +112,17 @@ protected int executeCountQuery(String query) { * query selects many columns and the mapper wants name-based access * instead of hardcoded indices. */ + // A column selected by the query but absent from the ResultSet means the Pinot schema has not + // caught up with the deployed code. Degrade that one field to null rather than failing the page. + protected String readString(ResultSet resultSet, int rowIndex, Map columnIndexMap, String column) { + Integer columnIndex = columnIndexMap.get(column); + if (columnIndex == null) { + log.debug("Column {} missing from Pinot result set - returning null", column); + return null; + } + return resultSet.getString(rowIndex, columnIndex); + } + protected Map buildColumnIndexMap(ResultSet resultSet) { Map columnIndexMap = new HashMap<>(); for (int i = 0; i < resultSet.getColumnCount(); i++) { diff --git a/openframe-data-pinot/src/main/java/com/openframe/data/pinot/repository/PinotClientLogRepository.java b/openframe-data-pinot/src/main/java/com/openframe/data/pinot/repository/PinotClientLogRepository.java index 37304ed6cc..d466890751 100644 --- a/openframe-data-pinot/src/main/java/com/openframe/data/pinot/repository/PinotClientLogRepository.java +++ b/openframe-data-pinot/src/main/java/com/openframe/data/pinot/repository/PinotClientLogRepository.java @@ -4,6 +4,7 @@ import com.openframe.data.pinot.model.OrganizationOption; import lombok.extern.slf4j.Slf4j; import org.apache.pinot.client.Connection; +import org.apache.pinot.client.ResultSet; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Repository; @@ -15,6 +16,8 @@ import java.util.Objects; import java.util.stream.Collectors; +import static org.springframework.util.StringUtils.hasText; + @Slf4j @Repository public class PinotClientLogRepository extends AbstractPinotRepository implements PinotLogRepository { @@ -178,23 +181,33 @@ public String getDefaultSortField() { return DEFAULT_SORT_COLUMN; } + // Most devices have no nickname: Pinot stores the schema default (empty string) for those rows, + // and the API contract is an absent nickname, not an empty one. + private String readNickname(ResultSet resultSet, int rowIndex, Map columnIndexMap) { + String nickname = readString(resultSet, rowIndex, columnIndexMap, "nickname"); + if (!hasText(nickname)) { + return null; + } + return nickname; + } + private List executeLogQuery(String query) { return executeQuery(query, resultSet -> { Map columnIndexMap = buildColumnIndexMap(resultSet); return rowIndex -> { LogProjection projection = new LogProjection(); - projection.toolEventId = resultSet.getString(rowIndex, columnIndexMap.get("toolEventId")); - projection.ingestDay = resultSet.getString(rowIndex, columnIndexMap.get("ingestDay")); - projection.toolType = resultSet.getString(rowIndex, columnIndexMap.get("toolType")); - projection.eventType = resultSet.getString(rowIndex, columnIndexMap.get("eventType")); - projection.severity = resultSet.getString(rowIndex, columnIndexMap.get("severity")); - projection.userId = resultSet.getString(rowIndex, columnIndexMap.get("userId")); - projection.deviceId = resultSet.getString(rowIndex, columnIndexMap.get("deviceId")); - projection.hostname = resultSet.getString(rowIndex, columnIndexMap.get("hostname")); - projection.nickname = resultSet.getString(rowIndex, columnIndexMap.get("nickname")); - projection.organizationId = resultSet.getString(rowIndex, columnIndexMap.get("organizationId")); - projection.organizationName = resultSet.getString(rowIndex, columnIndexMap.get("organizationName")); - projection.summary = resultSet.getString(rowIndex, columnIndexMap.get("summary")); + projection.toolEventId = readString(resultSet, rowIndex, columnIndexMap, "toolEventId"); + projection.ingestDay = readString(resultSet, rowIndex, columnIndexMap, "ingestDay"); + projection.toolType = readString(resultSet, rowIndex, columnIndexMap, "toolType"); + projection.eventType = readString(resultSet, rowIndex, columnIndexMap, "eventType"); + projection.severity = readString(resultSet, rowIndex, columnIndexMap, "severity"); + projection.userId = readString(resultSet, rowIndex, columnIndexMap, "userId"); + projection.deviceId = readString(resultSet, rowIndex, columnIndexMap, "deviceId"); + projection.hostname = readString(resultSet, rowIndex, columnIndexMap, "hostname"); + projection.nickname = readNickname(resultSet, rowIndex, columnIndexMap); + projection.organizationId = readString(resultSet, rowIndex, columnIndexMap, "organizationId"); + projection.organizationName = readString(resultSet, rowIndex, columnIndexMap, "organizationName"); + projection.summary = readString(resultSet, rowIndex, columnIndexMap, "summary"); projection.eventTimestamp = Instant.ofEpochMilli(resultSet.getLong(rowIndex, columnIndexMap.get("eventTimestamp"))); return projection; }; diff --git a/openframe-data-pinot/src/main/java/com/openframe/data/pinot/repository/PinotQueryBuilder.java b/openframe-data-pinot/src/main/java/com/openframe/data/pinot/repository/PinotQueryBuilder.java index f79e4b79d8..9423ab4016 100644 --- a/openframe-data-pinot/src/main/java/com/openframe/data/pinot/repository/PinotQueryBuilder.java +++ b/openframe-data-pinot/src/main/java/com/openframe/data/pinot/repository/PinotQueryBuilder.java @@ -439,6 +439,8 @@ public PinotQueryBuilder whereRelevanceLogSearch(String searchTerm) { relevanceConditions.add(TEXT_MATCH_FUNCTION + "(userId, '" + escapeSqlValue(processedSearchTerm) + "')"); + relevanceConditions.add(TEXT_MATCH_FUNCTION + "(nickname, '" + escapeSqlValue(processedSearchTerm) + "')"); + String relevanceCondition = "(" + String.join(SQL_OR, relevanceConditions) + ")"; whereConditions.add(relevanceCondition); } diff --git a/openframe-data-pinot/src/test/java/com/openframe/data/pinot/repository/PinotClientLogRepositoryTest.java b/openframe-data-pinot/src/test/java/com/openframe/data/pinot/repository/PinotClientLogRepositoryTest.java index adc77c609f..8ca5c8c7e2 100644 --- a/openframe-data-pinot/src/test/java/com/openframe/data/pinot/repository/PinotClientLogRepositoryTest.java +++ b/openframe-data-pinot/src/test/java/com/openframe/data/pinot/repository/PinotClientLogRepositoryTest.java @@ -21,6 +21,7 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.anyString; @@ -223,6 +224,91 @@ void mapsLogProjection() { assertNotNull(p.eventTimestamp); } + @Test + @DisplayName("a column missing from the result set degrades to null instead of failing the whole page") + void missingColumnDoesNotBreakMapping() { + // Pinot schema not yet caught up with the deployed code: nickname is selected but absent. + when(pinotConnection.execute(anyString())).thenReturn(resultSetGroup); + when(resultSetGroup.getResultSet(0)).thenReturn(resultSet); + when(resultSet.getRowCount()).thenReturn(1); + when(resultSet.getColumnCount()).thenReturn(12); + when(resultSet.getColumnName(0)).thenReturn("toolEventId"); + when(resultSet.getColumnName(1)).thenReturn("ingestDay"); + when(resultSet.getColumnName(2)).thenReturn("toolType"); + when(resultSet.getColumnName(3)).thenReturn("eventType"); + when(resultSet.getColumnName(4)).thenReturn("severity"); + when(resultSet.getColumnName(5)).thenReturn("userId"); + when(resultSet.getColumnName(6)).thenReturn("deviceId"); + when(resultSet.getColumnName(7)).thenReturn("hostname"); + when(resultSet.getColumnName(8)).thenReturn("organizationId"); + when(resultSet.getColumnName(9)).thenReturn("organizationName"); + when(resultSet.getColumnName(10)).thenReturn("summary"); + when(resultSet.getColumnName(11)).thenReturn("eventTimestamp"); + + when(resultSet.getString(0, 0)).thenReturn("toolEventId-v"); + when(resultSet.getString(0, 1)).thenReturn("ingestDay-v"); + when(resultSet.getString(0, 2)).thenReturn("toolType-v"); + when(resultSet.getString(0, 3)).thenReturn("eventType-v"); + when(resultSet.getString(0, 4)).thenReturn("severity-v"); + when(resultSet.getString(0, 5)).thenReturn("userId-v"); + when(resultSet.getString(0, 6)).thenReturn("deviceId-v"); + when(resultSet.getString(0, 7)).thenReturn("host-1"); + when(resultSet.getString(0, 8)).thenReturn("organizationId-v"); + when(resultSet.getString(0, 9)).thenReturn("organizationName-v"); + when(resultSet.getString(0, 10)).thenReturn("summary-v"); + when(resultSet.getLong(0, 11)).thenReturn(1700000000000L); + + List result = repository.findLogs(TENANT_ID, START, END, null, null, List.of(), List.of(), + List.of(), List.of(), null, null, 100, "eventTimestamp", "DESC"); + + assertEquals(1, result.size()); + LogProjection p = result.get(0); + assertEquals("host-1", p.hostname); + assertNull(p.nickname); + } + + @Test + @DisplayName("an empty nickname (the Pinot schema default) maps to null, not an empty string") + void emptyNicknameMapsToNull() { + when(pinotConnection.execute(anyString())).thenReturn(resultSetGroup); + when(resultSetGroup.getResultSet(0)).thenReturn(resultSet); + when(resultSet.getRowCount()).thenReturn(1); + when(resultSet.getColumnCount()).thenReturn(13); + when(resultSet.getColumnName(0)).thenReturn("toolEventId"); + when(resultSet.getColumnName(1)).thenReturn("ingestDay"); + when(resultSet.getColumnName(2)).thenReturn("toolType"); + when(resultSet.getColumnName(3)).thenReturn("eventType"); + when(resultSet.getColumnName(4)).thenReturn("severity"); + when(resultSet.getColumnName(5)).thenReturn("userId"); + when(resultSet.getColumnName(6)).thenReturn("deviceId"); + when(resultSet.getColumnName(7)).thenReturn("hostname"); + when(resultSet.getColumnName(8)).thenReturn("nickname"); + when(resultSet.getColumnName(9)).thenReturn("organizationId"); + when(resultSet.getColumnName(10)).thenReturn("organizationName"); + when(resultSet.getColumnName(11)).thenReturn("summary"); + when(resultSet.getColumnName(12)).thenReturn("eventTimestamp"); + + when(resultSet.getString(0, 0)).thenReturn("toolEventId-v"); + when(resultSet.getString(0, 1)).thenReturn("ingestDay-v"); + when(resultSet.getString(0, 2)).thenReturn("toolType-v"); + when(resultSet.getString(0, 3)).thenReturn("eventType-v"); + when(resultSet.getString(0, 4)).thenReturn("severity-v"); + when(resultSet.getString(0, 5)).thenReturn("userId-v"); + when(resultSet.getString(0, 6)).thenReturn("deviceId-v"); + when(resultSet.getString(0, 7)).thenReturn("host-1"); + when(resultSet.getString(0, 8)).thenReturn(""); + when(resultSet.getString(0, 9)).thenReturn("organizationId-v"); + when(resultSet.getString(0, 10)).thenReturn("organizationName-v"); + when(resultSet.getString(0, 11)).thenReturn("summary-v"); + when(resultSet.getLong(0, 12)).thenReturn(1700000000000L); + + List result = repository.findLogs(TENANT_ID, START, END, null, null, List.of(), List.of(), + List.of(), List.of(), null, null, 100, "eventTimestamp", "DESC"); + + assertEquals("host-1", result.get(0).hostname); + assertNull(result.get(0).nickname); + } + @Test @DisplayName("getOrganizationOptions filters null/empty organizationIds") void organizationOptionsFilterBlanks() { diff --git a/openframe-pinot-initializer/src/main/java/com/openframe/management/config/pinot/PinotConfigInitializer.java b/openframe-pinot-initializer/src/main/java/com/openframe/management/config/pinot/PinotConfigInitializer.java index 1cf532f3fd..dfe2c5e4d7 100644 --- a/openframe-pinot-initializer/src/main/java/com/openframe/management/config/pinot/PinotConfigInitializer.java +++ b/openframe-pinot-initializer/src/main/java/com/openframe/management/config/pinot/PinotConfigInitializer.java @@ -26,6 +26,7 @@ import org.springframework.web.client.ResourceAccessException; import org.springframework.web.client.RestTemplate; +import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; @@ -108,6 +109,8 @@ private void deployPinotConfig(PinotConfig config) { } + deployWithRetry(() -> reloadSegments(realtimeTableConfig), "segment reload for " + config.getName()); + log.info("Successfully deployed Pinot configuration for: {}", config.getName()); } catch (Exception e) { @@ -191,10 +194,7 @@ private void deploySchema(String schemaConfig) { private void deployTableConfig(String tableConfig, String configName) { try { - JsonNode tableConfigJson = objectMapper.readTree(tableConfig); - String baseTableName = tableConfigJson.get("tableName").asText(); - String tableType = tableConfigJson.get("tableType").asText(); - String tableName = baseTableName + ("REALTIME".equalsIgnoreCase(tableType) ? "_REALTIME" : "_OFFLINE"); + String tableName = resolveTableNameWithType(tableConfig); String updateUrl = String.format("http://%s/tables/%s", pinotControllerUrl, tableName); String createUrl = String.format("http://%s/tables", pinotControllerUrl); @@ -229,6 +229,44 @@ private void deployTableConfig(String tableConfig, String configName) { } } + // A column added to the schema only materialises in already-persisted segments when they are + // reloaded; without this, queries selecting the new column see it missing on historical data. + private void reloadSegments(String tableConfig) { + try { + String tableName = resolveTableNameWithType(tableConfig); + String url = String.format("http://%s/segments/%s/reload", pinotControllerUrl, tableName); + + HttpHeaders headers = createHeaders(); + HttpEntity request = new HttpEntity<>(null, headers); + ResponseEntity response = restTemplate.exchange(url, HttpMethod.POST, request, String.class); + + if (response.getStatusCode() == HttpStatus.OK) { + log.info("Triggered segment reload for {}: {}", tableName, response.getBody()); + } else { + log.error("Failed to trigger segment reload for {}. Status: {}", tableName, response.getStatusCode()); + throw new RuntimeException("Failed to trigger segment reload. Status: " + response.getStatusCode()); + } + + } catch (Exception e) { + log.error("Error triggering segment reload", e); + throw new RuntimeException("Failed to trigger segment reload", e); + } + } + + private String resolveTableNameWithType(String tableConfig) { + try { + JsonNode tableConfigJson = objectMapper.readTree(tableConfig); + String baseTableName = tableConfigJson.get("tableName").asText(); + String tableType = tableConfigJson.get("tableType").asText(); + if ("REALTIME".equalsIgnoreCase(tableType)) { + return baseTableName + "_REALTIME"; + } + return baseTableName + "_OFFLINE"; + } catch (JsonProcessingException e) { + throw new RuntimeException("Failed to read Pinot table configuration", e); + } + } + private HttpHeaders createHeaders() { HttpHeaders headers = new HttpHeaders(); headers.setContentType(MediaType.APPLICATION_JSON); diff --git a/openframe-pinot-initializer/src/main/resources/pinot/config/schema-logs.json b/openframe-pinot-initializer/src/main/resources/pinot/config/schema-logs.json index 4d4eb1ac92..eb100a979f 100644 --- a/openframe-pinot-initializer/src/main/resources/pinot/config/schema-logs.json +++ b/openframe-pinot-initializer/src/main/resources/pinot/config/schema-logs.json @@ -7,7 +7,7 @@ { "name": "userId", "dataType": "STRING", "singleValueField": true }, { "name": "deviceId", "dataType": "STRING", "singleValueField": true }, { "name": "hostname", "dataType": "STRING", "singleValueField": true }, - { "name": "nickname", "dataType": "STRING", "singleValueField": true }, + { "name": "nickname", "dataType": "STRING", "singleValueField": true, "defaultNullValue": "" }, { "name": "organizationId", "dataType": "STRING", "singleValueField": true }, { "name": "organizationName", "dataType": "STRING", "singleValueField": true }, { "name": "ingestDay", "dataType": "STRING", "singleValueField": true }, diff --git a/openframe-stream-service-core/src/test/java/com/openframe/stream/handler/LogEventNicknamePropagationTest.java b/openframe-stream-service-core/src/test/java/com/openframe/stream/handler/LogEventNicknamePropagationTest.java new file mode 100644 index 0000000000..fd77850c4a --- /dev/null +++ b/openframe-stream-service-core/src/test/java/com/openframe/stream/handler/LogEventNicknamePropagationTest.java @@ -0,0 +1,105 @@ +package com.openframe.stream.handler; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.openframe.data.cassandra.model.UnifiedLogEvent; +import com.openframe.data.cassandra.model.enums.UnifiedEventType; +import com.openframe.data.cassandra.repository.UnifiedLogEventRepository; +import com.openframe.data.model.enums.IntegratedToolType; +import com.openframe.kafka.model.IntegratedToolEvent; +import com.openframe.kafka.model.debezium.DebeziumMessage; +import com.openframe.kafka.producer.retry.OssTenantRetryingKafkaProducer; +import com.openframe.stream.model.fleet.debezium.DeserializedDebeziumMessage; +import com.openframe.stream.model.fleet.debezium.IntegratedToolEnrichedData; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.isNull; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; + +/** + * The two log sinks must agree: whatever enrichment resolved has to reach both the Cassandra + * detail row and the Kafka message that Pinot ingests, or the list and detail views disagree. + */ +class LogEventNicknamePropagationTest { + + private static final String TENANT_ID = "tenant-a"; + private static final String MACHINE_ID = "6d925893-702a-4223-b62f-2f80b927cbaa"; + private static final String HOSTNAME = "MBP-Oleksandr.lan"; + private static final String NICKNAME = "Reception iMac"; + + private static DeserializedDebeziumMessage message() { + DebeziumMessage.Payload payload = new DebeziumMessage.Payload<>(); + payload.setOperation("c"); + return DeserializedDebeziumMessage.builder() + .payload(payload) + .tenantId(TENANT_ID) + .toolEventId("evt-1") + .ingestDay("2026-07-24") + .integratedToolType(IntegratedToolType.FLEET) + .unifiedEventType(UnifiedEventType.LOGIN) + .eventTimestamp(1L) + .isVisible(true) + .build(); + } + + private static IntegratedToolEnrichedData enriched(String nickname) { + IntegratedToolEnrichedData enriched = new IntegratedToolEnrichedData(); + enriched.setTenantId(TENANT_ID); + enriched.setMachineId(MACHINE_ID); + enriched.setHostname(HOSTNAME); + enriched.setNickname(nickname); + return enriched; + } + + @Test + @DisplayName("Cassandra sink: hostname and nickname are both written to unified_logs") + void cassandraHandlerWritesBothNameFields() { + UnifiedLogEventRepository repository = mock(UnifiedLogEventRepository.class); + DebeziumCassandraMessageHandler handler = new DebeziumCassandraMessageHandler( + repository, new ObjectMapper(), new TenantIdRequiredDebeziumEventValidator()); + + handler.handle(message(), enriched(NICKNAME)); + + ArgumentCaptor captor = ArgumentCaptor.forClass(UnifiedLogEvent.class); + verify(repository).save(captor.capture()); + UnifiedLogEvent saved = captor.getValue(); + assertThat(saved.getHostname()).isEqualTo(HOSTNAME); + assertThat(saved.getNickname()).isEqualTo(NICKNAME); + } + + @Test + @DisplayName("Cassandra sink: a machine with no nickname writes a null nickname, not the hostname") + void cassandraHandlerWritesNullNicknameWhenAbsent() { + UnifiedLogEventRepository repository = mock(UnifiedLogEventRepository.class); + DebeziumCassandraMessageHandler handler = new DebeziumCassandraMessageHandler( + repository, new ObjectMapper(), new TenantIdRequiredDebeziumEventValidator()); + + handler.handle(message(), enriched(null)); + + ArgumentCaptor captor = ArgumentCaptor.forClass(UnifiedLogEvent.class); + verify(repository).save(captor.capture()); + UnifiedLogEvent saved = captor.getValue(); + assertThat(saved.getHostname()).isEqualTo(HOSTNAME); + assertThat(saved.getNickname()).isNull(); + } + + @Test + @DisplayName("Kafka/Pinot sink: hostname and nickname both reach the published message") + void kafkaHandlerPublishesBothNameFields() { + OssTenantRetryingKafkaProducer producer = mock(OssTenantRetryingKafkaProducer.class); + TenantDebeziumKafkaMessageHandler handler = new TenantDebeziumKafkaMessageHandler( + producer, new ObjectMapper(), new TenantIdRequiredDebeziumEventValidator()); + + handler.handle(message(), enriched(NICKNAME)); + + ArgumentCaptor captor = ArgumentCaptor.forClass(IntegratedToolEvent.class); + verify(producer).publish(isNull(), anyString(), captor.capture()); + IntegratedToolEvent published = captor.getValue(); + assertThat(published.getHostname()).isEqualTo(HOSTNAME); + assertThat(published.getNickname()).isEqualTo(NICKNAME); + } +} diff --git a/openframe-stream-service-core/src/test/java/com/openframe/stream/service/IntegratedToolDataEnrichmentServiceTest.java b/openframe-stream-service-core/src/test/java/com/openframe/stream/service/IntegratedToolDataEnrichmentServiceTest.java new file mode 100644 index 0000000000..63f9e2109c --- /dev/null +++ b/openframe-stream-service-core/src/test/java/com/openframe/stream/service/IntegratedToolDataEnrichmentServiceTest.java @@ -0,0 +1,109 @@ +package com.openframe.stream.service; + +import com.openframe.data.model.enums.DataEnrichmentServiceType; +import com.openframe.data.model.enums.IntegratedToolType; +import com.openframe.data.model.redis.CachedMachineInfo; +import com.openframe.data.model.redis.CachedOrganizationInfo; +import com.openframe.data.repository.redis.MachineIdCacheService; +import com.openframe.data.service.TenantIdProvider; +import com.openframe.kafka.model.debezium.DebeziumMessage; +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.DisplayName; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.lenient; +import static org.mockito.Mockito.when; + +/** + * Enrichment for the external-tool path (Fleet MDM and MeshCentral) — by volume the main + * producer of log events. Nickname must travel alongside hostname, never replace it. + */ +@ExtendWith(MockitoExtension.class) +class IntegratedToolDataEnrichmentServiceTest { + + private static final String AGENT_ID = "fleet-host-4417"; + private static final String MACHINE_ID = "6d925893-702a-4223-b62f-2f80b927cbaa"; + private static final String ORG_ID = "e0521785-8fef-4ec3-b520-f99087ed988e"; + private static final String TENANT_ID = "tenant-1"; + private static final String HOSTNAME = "MBP-Oleksandr.lan"; + private static final String NICKNAME = "Reception iMac"; + + @Mock + private MachineIdCacheService machineIdCacheService; + @Mock + private TenantIdProvider tenantIdProvider; + + private IntegratedToolDataEnrichmentService service; + + @BeforeEach + void setUp() { + lenient().when(tenantIdProvider.getTenantId()).thenReturn(TENANT_ID); + // Tenant cluster mode — no ClusterTenantIdResolver bean. + service = new IntegratedToolDataEnrichmentService(machineIdCacheService, null, tenantIdProvider); + } + + private static DeserializedDebeziumMessage message(String agentId) { + DebeziumMessage.Payload payload = new DebeziumMessage.Payload<>(); + payload.setOperation("c"); + return DeserializedDebeziumMessage.builder() + .payload(payload) + .agentId(agentId) + .integratedToolType(IntegratedToolType.FLEET) + .build(); + } + + @Test + @DisplayName("getType: returns INTEGRATED_TOOLS_EVENTS — the routing key for Fleet and MeshCentral") + void getType_returnsIntegratedToolsEvents() { + assertThat(service.getType()).isEqualTo(DataEnrichmentServiceType.INTEGRATED_TOOLS_EVENTS); + } + + @Test + @DisplayName("getExtraParams: nickname is carried alongside hostname, not instead of it") + void getExtraParams_carriesNicknameAlongsideHostname() { + when(machineIdCacheService.getMachine(AGENT_ID)) + .thenReturn(new CachedMachineInfo(MACHINE_ID, HOSTNAME, NICKNAME, ORG_ID)); + when(machineIdCacheService.getOrganization(ORG_ID)) + .thenReturn(new CachedOrganizationInfo(ORG_ID, "Default")); + + IntegratedToolEnrichedData enriched = service.getExtraParams(message(AGENT_ID)); + + assertThat(enriched.getMachineId()).isEqualTo(MACHINE_ID); + assertThat(enriched.getHostname()).isEqualTo(HOSTNAME); + assertThat(enriched.getNickname()).isEqualTo(NICKNAME); + assertThat(enriched.getOrganizationId()).isEqualTo(ORG_ID); + assertThat(enriched.getTenantId()).isEqualTo(TENANT_ID); + } + + @Test + @DisplayName("getExtraParams: machine without a nickname leaves nickname null and keeps the hostname") + void getExtraParams_noNickname_leavesNicknameNull() { + when(machineIdCacheService.getMachine(AGENT_ID)) + .thenReturn(new CachedMachineInfo(MACHINE_ID, HOSTNAME, null, ORG_ID)); + when(machineIdCacheService.getOrganization(ORG_ID)) + .thenReturn(new CachedOrganizationInfo(ORG_ID, "Default")); + + IntegratedToolEnrichedData enriched = service.getExtraParams(message(AGENT_ID)); + + assertThat(enriched.getHostname()).isEqualTo(HOSTNAME); + assertThat(enriched.getNickname()).isNull(); + } + + @Test + @DisplayName("getExtraParams: unknown machine leaves both name fields null but still resolves the tenant") + void getExtraParams_unknownMachine_leavesNameFieldsNull() { + when(machineIdCacheService.getMachine(AGENT_ID)).thenReturn(null); + + IntegratedToolEnrichedData enriched = service.getExtraParams(message(AGENT_ID)); + + assertThat(enriched.getHostname()).isNull(); + assertThat(enriched.getNickname()).isNull(); + assertThat(enriched.getTenantId()).isEqualTo(TENANT_ID); + } +} diff --git a/openframe-stream-service-core/src/test/java/com/openframe/stream/service/ScriptExecutedEnrichmentIntegrationTest.java b/openframe-stream-service-core/src/test/java/com/openframe/stream/service/ScriptExecutedEnrichmentIntegrationTest.java index 02a8b9f963..46920ba7e4 100644 --- a/openframe-stream-service-core/src/test/java/com/openframe/stream/service/ScriptExecutedEnrichmentIntegrationTest.java +++ b/openframe-stream-service-core/src/test/java/com/openframe/stream/service/ScriptExecutedEnrichmentIntegrationTest.java @@ -50,6 +50,7 @@ class ScriptExecutedEnrichmentIntegrationTest { // ties to the bug it guards against. private static final String MACHINE_ID = "6d925893-702a-4223-b62f-2f80b927cbaa"; private static final String HOSTNAME = "MBP-Oleksandr.lan"; + private static final String NICKNAME = "Reception iMac"; private static final String ORG_ID = "e0521785-8fef-4ec3-b520-f99087ed988e"; private static final String ORG_NAME = "Default"; private static final String TENANT_ID = "tenant-1"; @@ -82,7 +83,7 @@ void scriptExecutedKafkaMessage_yieldsFullyPopulatedEnrichment() { // 3. Stub the Mongo/Redis-backed lookups; real cache layer is irrelevant here. when(machineIdCacheService.getMachineByMachineId(MACHINE_ID)) - .thenReturn(new CachedMachineInfo(MACHINE_ID, HOSTNAME, null, ORG_ID)); + .thenReturn(new CachedMachineInfo(MACHINE_ID, HOSTNAME, NICKNAME, ORG_ID)); when(machineIdCacheService.getOrganization(ORG_ID)) .thenReturn(new CachedOrganizationInfo(ORG_ID, ORG_NAME)); when(tenantIdProvider.getTenantId()).thenReturn(TENANT_ID); @@ -101,6 +102,9 @@ void scriptExecutedKafkaMessage_yieldsFullyPopulatedEnrichment() { assertThat(enriched.getHostname()) .as("hostname on the LogEvent") .isEqualTo(HOSTNAME); + assertThat(enriched.getNickname()) + .as("nickname on the LogEvent — carried alongside hostname, never replacing it") + .isEqualTo(NICKNAME); assertThat(enriched.getOrganizationId()) .as("organizationId on the LogEvent") .isEqualTo(ORG_ID); @@ -134,6 +138,7 @@ void commandExecutedKafkaMessage_yieldsFullyPopulatedEnrichment() { assertThat(enriched.getMachineId()).isEqualTo(MACHINE_ID); assertThat(enriched.getHostname()).isEqualTo(HOSTNAME); + assertThat(enriched.getNickname()).isNull(); assertThat(enriched.getOrganizationId()).isEqualTo(ORG_ID); assertThat(enriched.getOrganizationName()).isEqualTo(ORG_NAME); assertThat(enriched.getTenantId()).isEqualTo(TENANT_ID); From 33d0e0a79f2fae9f4c60331d67f672e102898a9b Mon Sep 17 00:00:00 2001 From: Hryhorii Chuhuievets Date: Mon, 14 Sep 2026 19:00:16 +0300 Subject: [PATCH 4/5] Add default values for hostname, organizationId and organizationName to avoid "null" value from fleet --- .../src/main/resources/pinot/config/schema-logs.json | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/openframe-pinot-initializer/src/main/resources/pinot/config/schema-logs.json b/openframe-pinot-initializer/src/main/resources/pinot/config/schema-logs.json index eb100a979f..1ac51b6bae 100644 --- a/openframe-pinot-initializer/src/main/resources/pinot/config/schema-logs.json +++ b/openframe-pinot-initializer/src/main/resources/pinot/config/schema-logs.json @@ -6,10 +6,10 @@ { "name": "toolEventId", "dataType": "STRING", "singleValueField": true }, { "name": "userId", "dataType": "STRING", "singleValueField": true }, { "name": "deviceId", "dataType": "STRING", "singleValueField": true }, - { "name": "hostname", "dataType": "STRING", "singleValueField": true }, + { "name": "hostname", "dataType": "STRING", "singleValueField": true, "defaultNullValue": "" }, { "name": "nickname", "dataType": "STRING", "singleValueField": true, "defaultNullValue": "" }, - { "name": "organizationId", "dataType": "STRING", "singleValueField": true }, - { "name": "organizationName", "dataType": "STRING", "singleValueField": true }, + { "name": "organizationId", "dataType": "STRING", "singleValueField": true, "defaultNullValue": "" }, + { "name": "organizationName", "dataType": "STRING", "singleValueField": true, "defaultNullValue": "" }, { "name": "ingestDay", "dataType": "STRING", "singleValueField": true }, { "name": "toolType", "dataType": "STRING", "singleValueField": true }, { "name": "eventType", "dataType": "STRING", "singleValueField": true }, From c084783666d73b9fda2cd6e6c61d99d21c00633d Mon Sep 17 00:00:00 2001 From: Hryhorii Chuhuievets Date: Thu, 17 Sep 2026 12:02:09 +0300 Subject: [PATCH 5/5] Add cache evict for Machine, once nickname is updated --- openframe-api-lib/pom.xml | 4 + .../api/event/DeviceNicknameUpdatedEvent.java | 15 +++ .../api/service/device/DeviceService.java | 4 + .../MachineCacheInvalidationPublisher.java | 32 +++++++ ...MachineCacheInvalidationPublisherTest.java | 47 ++++++++++ .../api/service/DeviceServiceTest.java | 32 ++++++- .../redis/MachineIdCacheService.java | 45 ++++++++- .../redis/MachineIdCacheServiceTest.java | 92 +++++++++++++++++++ .../MachineCacheInvalidationConfig.java | 27 ++++++ .../MachineCacheInvalidationListener.java | 27 ++++++ .../MachineCacheInvalidationListenerTest.java | 38 ++++++++ 11 files changed, 359 insertions(+), 4 deletions(-) create mode 100644 openframe-api-lib/src/main/java/com/openframe/api/event/DeviceNicknameUpdatedEvent.java create mode 100644 openframe-api-lib/src/main/java/com/openframe/api/service/device/MachineCacheInvalidationPublisher.java create mode 100644 openframe-api-lib/src/test/java/com/openframe/api/service/device/MachineCacheInvalidationPublisherTest.java create mode 100644 openframe-data-mongo-sync/src/test/java/com/openframe/data/repository/redis/MachineIdCacheServiceTest.java create mode 100644 openframe-stream-service-core/src/main/java/com/openframe/stream/config/MachineCacheInvalidationConfig.java create mode 100644 openframe-stream-service-core/src/main/java/com/openframe/stream/listener/MachineCacheInvalidationListener.java create mode 100644 openframe-stream-service-core/src/test/java/com/openframe/stream/listener/MachineCacheInvalidationListenerTest.java diff --git a/openframe-api-lib/pom.xml b/openframe-api-lib/pom.xml index 7818777252..a2300ff4ed 100644 --- a/openframe-api-lib/pom.xml +++ b/openframe-api-lib/pom.xml @@ -25,6 +25,10 @@ com.openframe.oss openframe-data-mongo-sync + + com.openframe.oss + openframe-data-redis + diff --git a/openframe-api-lib/src/main/java/com/openframe/api/event/DeviceNicknameUpdatedEvent.java b/openframe-api-lib/src/main/java/com/openframe/api/event/DeviceNicknameUpdatedEvent.java new file mode 100644 index 0000000000..116b8bad06 --- /dev/null +++ b/openframe-api-lib/src/main/java/com/openframe/api/event/DeviceNicknameUpdatedEvent.java @@ -0,0 +1,15 @@ +package com.openframe.api.event; + +import lombok.Getter; +import org.springframework.context.ApplicationEvent; + +@Getter +public class DeviceNicknameUpdatedEvent extends ApplicationEvent { + + private final String machineId; + + public DeviceNicknameUpdatedEvent(Object source, String machineId) { + super(source); + this.machineId = machineId; + } +} diff --git a/openframe-api-lib/src/main/java/com/openframe/api/service/device/DeviceService.java b/openframe-api-lib/src/main/java/com/openframe/api/service/device/DeviceService.java index 6626ccc0ff..eb5b3191f2 100644 --- a/openframe-api-lib/src/main/java/com/openframe/api/service/device/DeviceService.java +++ b/openframe-api-lib/src/main/java/com/openframe/api/service/device/DeviceService.java @@ -8,6 +8,7 @@ import com.openframe.api.dto.shared.PageInfo; import com.openframe.api.dto.shared.SortDirection; import com.openframe.api.dto.shared.SortInput; +import com.openframe.api.event.DeviceNicknameUpdatedEvent; import com.openframe.api.exception.DeviceNotFoundException; import com.openframe.api.mapper.DeviceFilterOptionMapper; import com.openframe.api.service.processor.DeviceStatusProcessor; @@ -32,6 +33,7 @@ import jakarta.validation.constraints.NotNull; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import org.springframework.context.ApplicationEventPublisher; import org.springframework.stereotype.Service; import org.springframework.validation.annotation.Validated; @@ -64,6 +66,7 @@ public class DeviceService { private final ScheduleScriptDeviceService scheduleScriptDeviceService; private final DeviceFilterOptionMapper deviceFilterOptionMapper; private final TenantIdProvider tenantIdProvider; + private final ApplicationEventPublisher eventPublisher; public Optional findByMachineId(@NotBlank String machineId) { log.debug("Finding machine by ID: {}", machineId); @@ -368,6 +371,7 @@ public Machine updateNickname(@NotBlank String machineId, String nickname) { MachineWriteResult result = machineWriter .update(machineId, machineUpdate().set(NICKNAME, normalizeNickname(nickname))) .orElseThrow(() -> new DeviceNotFoundException("Device not found: " + machineId)); + eventPublisher.publishEvent(new DeviceNicknameUpdatedEvent(this, machineId)); log.info("Device {} nickname updated", machineId); return result.after(); } diff --git a/openframe-api-lib/src/main/java/com/openframe/api/service/device/MachineCacheInvalidationPublisher.java b/openframe-api-lib/src/main/java/com/openframe/api/service/device/MachineCacheInvalidationPublisher.java new file mode 100644 index 0000000000..e5ea1ffe67 --- /dev/null +++ b/openframe-api-lib/src/main/java/com/openframe/api/service/device/MachineCacheInvalidationPublisher.java @@ -0,0 +1,32 @@ +package com.openframe.api.service.device; + +import com.openframe.api.event.DeviceNicknameUpdatedEvent; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.context.event.EventListener; +import org.springframework.data.redis.core.StringRedisTemplate; +import org.springframework.stereotype.Component; + +import static com.openframe.data.repository.redis.MachineIdCacheService.INVALIDATION_CHANNEL; + +@Component +@Slf4j +@RequiredArgsConstructor +@ConditionalOnProperty(name = "spring.redis.enabled", havingValue = "true") +public class MachineCacheInvalidationPublisher { + + private final StringRedisTemplate redisTemplate; + + @EventListener + public void onNicknameUpdated(DeviceNicknameUpdatedEvent event) { + String machineId = event.getMachineId(); + try { + redisTemplate.convertAndSend(INVALIDATION_CHANNEL, machineId); + log.info("Published machine cache invalidation: machineId={}", machineId); + } catch (Exception e) { + log.warn("Failed to publish machine cache invalidation, cached info expires by TTL: machineId={}", + machineId, e); + } + } +} diff --git a/openframe-api-lib/src/test/java/com/openframe/api/service/device/MachineCacheInvalidationPublisherTest.java b/openframe-api-lib/src/test/java/com/openframe/api/service/device/MachineCacheInvalidationPublisherTest.java new file mode 100644 index 0000000000..30360b561c --- /dev/null +++ b/openframe-api-lib/src/test/java/com/openframe/api/service/device/MachineCacheInvalidationPublisherTest.java @@ -0,0 +1,47 @@ +package com.openframe.api.service.device; + +import com.openframe.api.event.DeviceNicknameUpdatedEvent; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.InjectMocks; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.data.redis.RedisConnectionFailureException; +import org.springframework.data.redis.core.StringRedisTemplate; + +import static com.openframe.data.repository.redis.MachineIdCacheService.INVALIDATION_CHANNEL; +import static org.assertj.core.api.Assertions.assertThatCode; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.verify; + +@ExtendWith(MockitoExtension.class) +class MachineCacheInvalidationPublisherTest { + + private static final String MACHINE_ID = "6d925893-702a-4223-b62f-2f80b927cbaa"; + + @Mock + private StringRedisTemplate redisTemplate; + + @InjectMocks + private MachineCacheInvalidationPublisher publisher; + + @Test + @DisplayName("onNicknameUpdated: forwards the machineId on the global invalidation channel the stream services subscribe to") + void onNicknameUpdated_publishesMachineIdOnGlobalChannel() { + publisher.onNicknameUpdated(new DeviceNicknameUpdatedEvent(this, MACHINE_ID)); + + verify(redisTemplate).convertAndSend(INVALIDATION_CHANNEL, MACHINE_ID); + } + + @Test + @DisplayName("onNicknameUpdated: a Redis failure never propagates — the rename is already persisted and the cache TTL covers it") + void onNicknameUpdated_redisFailureDoesNotPropagate() { + doThrow(new RedisConnectionFailureException("redis down")) + .when(redisTemplate).convertAndSend(any(), any()); + + assertThatCode(() -> publisher.onNicknameUpdated(new DeviceNicknameUpdatedEvent(this, MACHINE_ID))) + .doesNotThrowAnyException(); + } +} diff --git a/openframe-api-service-core/src/test/java/com/openframe/api/service/DeviceServiceTest.java b/openframe-api-service-core/src/test/java/com/openframe/api/service/DeviceServiceTest.java index 27e1b88b60..96eb93a835 100644 --- a/openframe-api-service-core/src/test/java/com/openframe/api/service/DeviceServiceTest.java +++ b/openframe-api-service-core/src/test/java/com/openframe/api/service/DeviceServiceTest.java @@ -2,6 +2,7 @@ import com.openframe.api.dto.device.DeviceFilterCriteria; import com.openframe.api.dto.shared.CursorPaginationCriteria; +import com.openframe.api.event.DeviceNicknameUpdatedEvent; import com.openframe.api.exception.DeviceNotFoundException; import com.openframe.api.mapper.DeviceFilterOptionMapper; import com.openframe.api.service.device.DeviceService; @@ -25,6 +26,7 @@ import org.mockito.ArgumentCaptor; import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.context.ApplicationEventPublisher; import java.time.Instant; import java.util.List; @@ -59,11 +61,14 @@ class DeviceServiceTest { @Mock private DeviceFilterOptionMapper deviceFilterOptionMapper; @Mock private TenantIdProvider tenantIdProvider; + @Mock + private ApplicationEventPublisher eventPublisher; private DeviceService service() { DeviceService s = new DeviceService(machineRepository, deviceOnlineDispatchRepository, machineWriter, tagRepository, tagAssignmentRepository, - deviceStatusProcessor, scheduleScriptDeviceService, deviceFilterOptionMapper, tenantIdProvider); + deviceStatusProcessor, scheduleScriptDeviceService, deviceFilterOptionMapper, tenantIdProvider, + eventPublisher); lenient().when(tenantIdProvider.getTenantId()).thenReturn(TENANT_ID); lenient().when(machineRepository.countMachines(any(), any(MachineQueryFilter.class), any())).thenReturn(0L); lenient().when(machineRepository.findMachinesWithCursor(any(), any(MachineQueryFilter.class), any(), @@ -369,6 +374,31 @@ void updateNickname_doesNotSaveWholeDocument() { verify(machineWriter).update(eq("m1"), any(MachineUpdate.class)); } + @Test + @DisplayName("updateNickname: publishes the invalidation event for the renamed machine once the write succeeded") + void updateNickname_publishesInvalidationEventAfterWrite() { + DeviceService s = service(); + stubAtomicNicknameUpdate("m1"); + + s.updateNickname("m1", "Reception iMac"); + + ArgumentCaptor captor = ArgumentCaptor.forClass(DeviceNicknameUpdatedEvent.class); + verify(eventPublisher).publishEvent(captor.capture()); + assertThat(captor.getValue().getMachineId()).isEqualTo("m1"); + } + + @Test + @DisplayName("updateNickname: an unknown device fails before anything is published — nothing to invalidate") + void updateNickname_unknownDevice_doesNotPublish() { + DeviceService s = service(); + when(machineWriter.update(eq("m1"), any(MachineUpdate.class))).thenReturn(Optional.empty()); + + assertThatThrownBy(() -> s.updateNickname("m1", "Reception iMac")) + .isInstanceOf(DeviceNotFoundException.class); + + verifyNoInteractions(eventPublisher); + } + @Test @DisplayName("updateNickname: a blank value clears the nickname (stored as null)") void updateNickname_blankClears() { diff --git a/openframe-data-mongo-sync/src/main/java/com/openframe/data/repository/redis/MachineIdCacheService.java b/openframe-data-mongo-sync/src/main/java/com/openframe/data/repository/redis/MachineIdCacheService.java index 1acd710c9a..fc050a911f 100644 --- a/openframe-data-mongo-sync/src/main/java/com/openframe/data/repository/redis/MachineIdCacheService.java +++ b/openframe-data-mongo-sync/src/main/java/com/openframe/data/repository/redis/MachineIdCacheService.java @@ -10,9 +10,13 @@ import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.cache.Cache; +import org.springframework.cache.CacheManager; import org.springframework.cache.annotation.Cacheable; import org.springframework.stereotype.Service; +import java.util.List; + /** * Service for machine and organization cache operations using Spring Cache abstraction * Uses lightweight DTOs to avoid serialization issues and reduce cache size @@ -24,9 +28,17 @@ @ConditionalOnProperty(name = "openframe.machine-id.cache.enabled", havingValue = "true") public class MachineIdCacheService { + public static final String INVALIDATION_CHANNEL = "machine:cache:invalidate"; + + private static final String MACHINE_CACHE = "machineCache"; + private static final String TENANT_MACHINE_CACHE = "tenantMachineCache"; + private static final String MACHINE_BY_ID_CACHE = "machineByIdCache"; + private static final String KEY_SEPARATOR = ":"; + private final ToolConnectionRepository toolConnectionRepository; private final MachineRepository machineRepository; private final OrganizationRepository organizationRepository; + private final CacheManager cacheManager; /** * Get cached machine info from cache or database by agent ID @@ -35,7 +47,7 @@ public class MachineIdCacheService { * @param agentId the agent ID * @return the CachedMachineInfo object, or null if not found */ - @Cacheable(value = "machineCache", key = "#agentId", unless = "#result == null") + @Cacheable(value = MACHINE_CACHE, key = "#agentId", unless = "#result == null") public CachedMachineInfo getMachine(String agentId) { log.debug("Fetching machine info for agent: {}", agentId); try { @@ -56,7 +68,7 @@ public CachedMachineInfo getMachine(String agentId) { } } - @Cacheable(value = "tenantMachineCache", key = "#tenantId + ':' + #toolType + ':' + #agentId", unless = "#result == null") + @Cacheable(value = TENANT_MACHINE_CACHE, key = "#tenantId + ':' + #toolType + ':' + #agentId", unless = "#result == null") public CachedMachineInfo getMachine(String tenantId, ToolType toolType, String agentId) { log.debug("Fetching machine info for agent: {} (tenant: {}, tool: {})", agentId, tenantId, toolType); try { @@ -84,7 +96,7 @@ public CachedMachineInfo getMachine(String tenantId, ToolType toolType, String a * @param machineId openframe-native machineId * @return the {@link CachedMachineInfo}, or {@code null} if the machine is not found in the local store */ - @Cacheable(value = "machineByIdCache", key = "#machineId", unless = "#result == null") + @Cacheable(value = MACHINE_BY_ID_CACHE, key = "#machineId", unless = "#result == null") public CachedMachineInfo getMachineByMachineId(String machineId) { log.debug("Fetching machine info by machineId: {}", machineId); try { @@ -124,5 +136,32 @@ public CachedOrganizationInfo getOrganization(String organizationId) { return null; } } + + public void evictMachine(String machineId) { + evict(MACHINE_BY_ID_CACHE, machineId); + List connections = toolConnectionRepository.findByMachineId(machineId); + connections.forEach(this::evictConnectionEntries); + log.info("Evicted cached machine info: machineId={} toolConnections={}", machineId, connections.size()); + } + + private void evictConnectionEntries(ToolConnection connection) { + String agentToolId = connection.getAgentToolId(); + String tenantMachineKey = tenantMachineKey(connection); + evict(MACHINE_CACHE, agentToolId); + evict(TENANT_MACHINE_CACHE, tenantMachineKey); + } + + private String tenantMachineKey(ToolConnection connection) { + return connection.getTenantId() + KEY_SEPARATOR + connection.getToolType() + + KEY_SEPARATOR + connection.getAgentToolId(); + } + + private void evict(String cacheName, String key) { + Cache cache = cacheManager.getCache(cacheName); + if (cache == null) { + return; + } + cache.evict(key); + } } diff --git a/openframe-data-mongo-sync/src/test/java/com/openframe/data/repository/redis/MachineIdCacheServiceTest.java b/openframe-data-mongo-sync/src/test/java/com/openframe/data/repository/redis/MachineIdCacheServiceTest.java new file mode 100644 index 0000000000..97f048c6d0 --- /dev/null +++ b/openframe-data-mongo-sync/src/test/java/com/openframe/data/repository/redis/MachineIdCacheServiceTest.java @@ -0,0 +1,92 @@ +package com.openframe.data.repository.redis; + +import com.openframe.data.document.tool.ToolConnection; +import com.openframe.data.document.tool.ToolType; +import com.openframe.data.repository.device.MachineRepository; +import com.openframe.data.repository.organization.OrganizationRepository; +import com.openframe.data.repository.tool.ToolConnectionRepository; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.cache.Cache; +import org.springframework.cache.CacheManager; + +import java.util.List; + +import static org.assertj.core.api.Assertions.assertThatCode; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoInteractions; +import static org.mockito.Mockito.when; + +@ExtendWith(MockitoExtension.class) +class MachineIdCacheServiceTest { + + private static final String MACHINE_ID = "6d925893-702a-4223-b62f-2f80b927cbaa"; + private static final String TENANT_ID = "tenant-1"; + private static final String AGENT_TOOL_ID = "node-42"; + + @Mock private ToolConnectionRepository toolConnectionRepository; + @Mock private MachineRepository machineRepository; + @Mock private OrganizationRepository organizationRepository; + @Mock private CacheManager cacheManager; + @Mock private Cache machineByIdCache; + @Mock private Cache machineCache; + @Mock private Cache tenantMachineCache; + + private MachineIdCacheService service; + + @BeforeEach + void setUp() { + service = new MachineIdCacheService(toolConnectionRepository, machineRepository, organizationRepository, + cacheManager); + } + + private static ToolConnection meshCentralConnection() { + ToolConnection connection = new ToolConnection(); + connection.setTenantId(TENANT_ID); + connection.setMachineId(MACHINE_ID); + connection.setToolType(ToolType.MESHCENTRAL); + connection.setAgentToolId(AGENT_TOOL_ID); + return connection; + } + + @Test + @DisplayName("evictMachine: evicts every key the enrichment reads through — by machineId, by agent, and the tenant-scoped agent key") + void evictMachine_evictsEveryKeyTheEnrichmentReadsThrough() { + when(cacheManager.getCache("machineByIdCache")).thenReturn(machineByIdCache); + when(cacheManager.getCache("machineCache")).thenReturn(machineCache); + when(cacheManager.getCache("tenantMachineCache")).thenReturn(tenantMachineCache); + when(toolConnectionRepository.findByMachineId(MACHINE_ID)).thenReturn(List.of(meshCentralConnection())); + + service.evictMachine(MACHINE_ID); + + verify(machineByIdCache).evict(MACHINE_ID); + verify(machineCache).evict(AGENT_TOOL_ID); + // Same shape as the @Cacheable SpEL key: tenantId + ':' + toolType + ':' + agentId + verify(tenantMachineCache).evict("tenant-1:MESHCENTRAL:node-42"); + } + + @Test + @DisplayName("evictMachine: a machine with no tool connections only evicts the machineId entry") + void evictMachine_noToolConnections_onlyEvictsById() { + when(cacheManager.getCache("machineByIdCache")).thenReturn(machineByIdCache); + when(toolConnectionRepository.findByMachineId(MACHINE_ID)).thenReturn(List.of()); + + service.evictMachine(MACHINE_ID); + + verify(machineByIdCache).evict(MACHINE_ID); + verifyNoInteractions(machineCache, tenantMachineCache); + } + + @Test + @DisplayName("evictMachine: a cache that is not configured is skipped instead of failing the invalidation") + void evictMachine_toleratesUnconfiguredCache() { + when(cacheManager.getCache("machineByIdCache")).thenReturn(null); + when(toolConnectionRepository.findByMachineId(MACHINE_ID)).thenReturn(List.of()); + + assertThatCode(() -> service.evictMachine(MACHINE_ID)).doesNotThrowAnyException(); + } +} diff --git a/openframe-stream-service-core/src/main/java/com/openframe/stream/config/MachineCacheInvalidationConfig.java b/openframe-stream-service-core/src/main/java/com/openframe/stream/config/MachineCacheInvalidationConfig.java new file mode 100644 index 0000000000..0e855caf34 --- /dev/null +++ b/openframe-stream-service-core/src/main/java/com/openframe/stream/config/MachineCacheInvalidationConfig.java @@ -0,0 +1,27 @@ +package com.openframe.stream.config; + +import com.openframe.stream.listener.MachineCacheInvalidationListener; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.data.redis.connection.RedisConnectionFactory; +import org.springframework.data.redis.listener.ChannelTopic; +import org.springframework.data.redis.listener.RedisMessageListenerContainer; + +import static com.openframe.data.repository.redis.MachineIdCacheService.INVALIDATION_CHANNEL; + +@Configuration +@ConditionalOnProperty(name = "openframe.machine-id.cache.enabled", havingValue = "true") +public class MachineCacheInvalidationConfig { + + @Bean + public RedisMessageListenerContainer machineCacheInvalidationListenerContainer( + RedisConnectionFactory connectionFactory, + MachineCacheInvalidationListener listener) { + RedisMessageListenerContainer container = new RedisMessageListenerContainer(); + ChannelTopic topic = new ChannelTopic(INVALIDATION_CHANNEL); + container.setConnectionFactory(connectionFactory); + container.addMessageListener(listener, topic); + return container; + } +} diff --git a/openframe-stream-service-core/src/main/java/com/openframe/stream/listener/MachineCacheInvalidationListener.java b/openframe-stream-service-core/src/main/java/com/openframe/stream/listener/MachineCacheInvalidationListener.java new file mode 100644 index 0000000000..d3e69f40c9 --- /dev/null +++ b/openframe-stream-service-core/src/main/java/com/openframe/stream/listener/MachineCacheInvalidationListener.java @@ -0,0 +1,27 @@ +package com.openframe.stream.listener; + +import com.openframe.data.repository.redis.MachineIdCacheService; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.data.redis.connection.Message; +import org.springframework.data.redis.connection.MessageListener; +import org.springframework.stereotype.Component; + +import java.nio.charset.StandardCharsets; + +@Component +@Slf4j +@RequiredArgsConstructor +@ConditionalOnProperty(name = "openframe.machine-id.cache.enabled", havingValue = "true") +public class MachineCacheInvalidationListener implements MessageListener { + + private final MachineIdCacheService machineIdCacheService; + + @Override + public void onMessage(Message message, byte[] pattern) { + String machineId = new String(message.getBody(), StandardCharsets.UTF_8); + log.info("Received machine cache invalidation: machineId={}", machineId); + machineIdCacheService.evictMachine(machineId); + } +} diff --git a/openframe-stream-service-core/src/test/java/com/openframe/stream/listener/MachineCacheInvalidationListenerTest.java b/openframe-stream-service-core/src/test/java/com/openframe/stream/listener/MachineCacheInvalidationListenerTest.java new file mode 100644 index 0000000000..5bbc370343 --- /dev/null +++ b/openframe-stream-service-core/src/test/java/com/openframe/stream/listener/MachineCacheInvalidationListenerTest.java @@ -0,0 +1,38 @@ +package com.openframe.stream.listener; + +import com.openframe.data.repository.redis.MachineIdCacheService; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.InjectMocks; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.data.redis.connection.DefaultMessage; + +import java.nio.charset.StandardCharsets; + +import static com.openframe.data.repository.redis.MachineIdCacheService.INVALIDATION_CHANNEL; +import static org.mockito.Mockito.verify; + +@ExtendWith(MockitoExtension.class) +class MachineCacheInvalidationListenerTest { + + private static final String MACHINE_ID = "6d925893-702a-4223-b62f-2f80b927cbaa"; + + @Mock + private MachineIdCacheService machineIdCacheService; + + @InjectMocks + private MachineCacheInvalidationListener listener; + + @Test + @DisplayName("onMessage: decodes the machineId from the channel payload and evicts its cached info") + void onMessage_evictsMachineFromPayload() { + byte[] channel = INVALIDATION_CHANNEL.getBytes(StandardCharsets.UTF_8); + byte[] body = MACHINE_ID.getBytes(StandardCharsets.UTF_8); + + listener.onMessage(new DefaultMessage(channel, body), null); + + verify(machineIdCacheService).evictMachine(MACHINE_ID); + } +}