Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -33,8 +33,12 @@ public class KafkaStreamsConfig {
* Bootstrap servers for Kafka Streams.
* Tenant deployments set {@code spring.oss-tenant.kafka.bootstrap-servers};
* shared/SaaS deployments set {@code spring.saas.kafka.bootstrap-servers}.
* {@code openframe.stream.kafka-streams.bootstrap-servers} overrides both — needed when a
* deployment has BOTH clusters configured but the streams topology must run against a
* specific one (e.g. the shared cluster's Fleet activity join reads the Debezium raw topics
* on the shared Kafka while {@code spring.oss-tenant.kafka} points at the tenant cluster).
*/
@Value("${spring.oss-tenant.kafka.bootstrap-servers:${spring.saas.kafka.bootstrap-servers:}}")
@Value("${openframe.stream.kafka-streams.bootstrap-servers:${spring.oss-tenant.kafka.bootstrap-servers:${spring.saas.kafka.bootstrap-servers:}}}")
private String bootstrapServers;

@Value("${spring.application.name}")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,13 @@ protected Optional<String> getAgentId(JsonNode after) {
return parseStringField(after, FIELD_AGENT_ID);
}

@Override
protected Optional<String> getTenantId(JsonNode after) {
// Stamped on activity_past by the Fleet fork and carried through the activity/host
// Kafka Streams join (Activity.teamId maps the team_id column).
return extractFleetTeamId(after);
}

@Override
protected Optional<String> getSourceEventType(JsonNode after) {
// Fleet MDM stores the event type in the "activity_type" column
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,12 +3,15 @@
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.openframe.data.model.enums.IntegratedToolType;
import com.openframe.data.model.enums.MessageType;
import com.openframe.sdk.fleetmdm.model.Policy;
import com.openframe.stream.mapping.FleetActivityTypeMapping;
import com.openframe.stream.service.ClusterTenantIdResolver;
import com.openframe.stream.service.FleetMdmCacheService;
import com.openframe.stream.util.TimestampParser;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;

import java.util.List;
Expand All @@ -30,10 +33,23 @@ public class FleetPolicyActivityDeserializer extends IntegratedToolEventDeserial
);

private final FleetMdmCacheService fleetMdmCacheService;
private final ClusterTenantIdResolver clusterTenantIdResolver;

protected FleetPolicyActivityDeserializer(ObjectMapper mapper, FleetMdmCacheService fleetMdmCacheService) {
protected FleetPolicyActivityDeserializer(ObjectMapper mapper, FleetMdmCacheService fleetMdmCacheService,
@Autowired(required = false) ClusterTenantIdResolver clusterTenantIdResolver) {
super(mapper, List.of(), List.of());
this.fleetMdmCacheService = fleetMdmCacheService;
this.clusterTenantIdResolver = clusterTenantIdResolver;
}

/** Shared cluster only: event tenant for the Fleet API lookup (see FleetQueryResultEventDeserializer). */
private String eventTenantId(JsonNode afterField) {
if (clusterTenantIdResolver == null) {
return null;
}
return extractFleetTeamId(afterField)
.map(teamId -> clusterTenantIdResolver.resolveTenantId(IntegratedToolType.FLEET, teamId))
.orElse(null);
}

@Override
Expand All @@ -46,6 +62,14 @@ protected Optional<String> getAgentId(JsonNode after) {
return parseStringField(after, FIELD_AGENT_ID);
}

@Override
protected Optional<String> getTenantId(JsonNode after) {
// Stamped on activity_past by the Fleet fork and carried through the activity/host
// Kafka Streams join (Activity.teamId maps the team_id column) — without this the
// shared cluster would drop every policy-CRUD activity as tenant-unresolved.
return extractFleetTeamId(after);
}

@Override
protected Optional<String> getSourceEventType(JsonNode after) {
return parseStringField(after, FIELD_ACTIVITY_TYPE);
Expand Down Expand Up @@ -133,13 +157,14 @@ private Optional<Policy> getPolicyInfo(JsonNode after) {
Optional<Long> policyIdOpt = extractPolicyId(after);

// Evict cache on policy mutation events so subsequent lookups get fresh data
String eventTenantId = eventTenantId(after);
policyIdOpt.ifPresent(policyId ->
getSourceEventType(after)
.filter(POLICY_MUTATION_TYPES::contains)
.ifPresent(type -> fleetMdmCacheService.evictPolicyCache(policyId))
.ifPresent(type -> fleetMdmCacheService.evictPolicyCache(policyId, eventTenantId))
);

return policyIdOpt.flatMap(fleetMdmCacheService::getPolicyById);
return policyIdOpt.flatMap(policyId -> fleetMdmCacheService.getPolicyById(policyId, eventTenantId));
}

private Optional<Long> extractPolicyId(JsonNode after) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,11 +3,14 @@
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.openframe.data.model.enums.IntegratedToolType;
import com.openframe.data.model.enums.MessageType;
import com.openframe.sdk.fleetmdm.model.Policy;
import com.openframe.stream.service.ClusterTenantIdResolver;
import com.openframe.stream.service.FleetMdmCacheService;
import com.openframe.stream.util.TimestampParser;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;

import java.util.List;
Expand All @@ -21,10 +24,23 @@
public class FleetPolicyMembershipEventDeserializer extends IntegratedToolEventDeserializer {

private final FleetMdmCacheService fleetMdmCacheService;
private final ClusterTenantIdResolver clusterTenantIdResolver;

protected FleetPolicyMembershipEventDeserializer(ObjectMapper mapper, FleetMdmCacheService fleetMdmCacheService) {
protected FleetPolicyMembershipEventDeserializer(ObjectMapper mapper, FleetMdmCacheService fleetMdmCacheService,
@Autowired(required = false) ClusterTenantIdResolver clusterTenantIdResolver) {
super(mapper, List.of(), List.of());
this.fleetMdmCacheService = fleetMdmCacheService;
this.clusterTenantIdResolver = clusterTenantIdResolver;
}

/** Shared cluster only: event tenant for the Fleet API lookup (see FleetQueryResultEventDeserializer). */
private String eventTenantId(JsonNode afterField) {
if (clusterTenantIdResolver == null) {
return null;
}
return extractFleetTeamId(afterField)
.map(teamId -> clusterTenantIdResolver.resolveTenantId(IntegratedToolType.FLEET, teamId))
.orElse(null);
}

@Override
Expand All @@ -38,6 +54,11 @@ protected Optional<String> getAgentId(JsonNode afterField) {
.map(JsonNode::asText);
}

@Override
protected Optional<String> getTenantId(JsonNode afterField) {
return extractFleetTeamId(afterField);
}

@Override
protected Optional<String> getSourceEventType(JsonNode afterField) {
JsonNode passesNode = afterField.get("passes");
Expand Down Expand Up @@ -128,6 +149,6 @@ private Optional<Policy> getPolicyInfo(JsonNode afterField) {
return Optional.ofNullable(afterField.get("policy_id"))
.filter(node -> !node.isNull())
.map(JsonNode::asLong)
.flatMap(fleetMdmCacheService::getPolicyById);
.flatMap(policyId -> fleetMdmCacheService.getPolicyById(policyId, eventTenantId(afterField)));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -3,10 +3,13 @@
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.openframe.data.model.enums.IntegratedToolType;
import com.openframe.data.model.enums.MessageType;
import com.openframe.sdk.fleetmdm.model.Query;
import com.openframe.stream.service.ClusterTenantIdResolver;
import com.openframe.stream.service.FleetMdmCacheService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;

import java.util.List;
Expand All @@ -19,10 +22,28 @@
public class FleetQueryResultEventDeserializer extends IntegratedToolEventDeserializer {

private final FleetMdmCacheService fleetMdmCacheService;
private final ClusterTenantIdResolver clusterTenantIdResolver;

protected FleetQueryResultEventDeserializer(ObjectMapper mapper, FleetMdmCacheService fleetMdmCacheService) {
protected FleetQueryResultEventDeserializer(ObjectMapper mapper, FleetMdmCacheService fleetMdmCacheService,
@Autowired(required = false) ClusterTenantIdResolver clusterTenantIdResolver) {
super(mapper, List.of(), List.of());
this.fleetMdmCacheService = fleetMdmCacheService;
this.clusterTenantIdResolver = clusterTenantIdResolver;
}

/**
* Shared cluster only: resolve the event row's stamped team_id to its tenant so the Fleet
* API lookup carries the right X-Tenant-Id (the shared Fleet's fences 404 a query fetched
* under the wrong tenant). Per-tenant clusters (no resolver bean) return null and the
* deployment client is used.
*/
private String eventTenantId(JsonNode afterField) {
if (clusterTenantIdResolver == null) {
return null;
}
return extractFleetTeamId(afterField)
.map(teamId -> clusterTenantIdResolver.resolveTenantId(IntegratedToolType.FLEET, teamId))
.orElse(null);
}

@Override
Expand All @@ -37,6 +58,11 @@ protected Optional<String> getAgentId(JsonNode afterField) {
.map(JsonNode::asText);
}

@Override
protected Optional<String> getTenantId(JsonNode afterField) {
return extractFleetTeamId(afterField);
}

@Override
protected Optional<String> getSourceEventType(JsonNode afterField) {
return Optional.of(EXECUTE_SCHEDULED_QUERY);
Expand Down Expand Up @@ -182,7 +208,7 @@ private Query getQueryInfo(JsonNode afterField) {
log.debug("Resolving query info for query_id: {}, host_id: {}", queryId,
afterField.has("host_id") ? afterField.get("host_id").asText() : "unknown");

Query query = fleetMdmCacheService.getQueryById(queryId);
Query query = fleetMdmCacheService.getQueryById(queryId, eventTenantId(afterField));
Comment thread
coderabbitai[bot] marked this conversation as resolved.

if (query == null) {
log.warn("Failed to resolve query name for query_id: {}. Fleet MDM client may not be initialized or query may have been deleted.", queryId);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -155,6 +155,24 @@ protected Optional<String> getTenantId(JsonNode afterField) {
return Optional.empty();
}

/**
* Tenant discriminator for Fleet CDC rows under shared-DB multitenancy: the {@code team_id}
* column stamped by the Fleet fork. Fleet deserializers return it from
* {@link #getTenantId(JsonNode)}; the shared cluster's {@code ClusterTenantIdResolver}
* maps it to the canonical tenant. Empty for rows written with the flag off (in per-tenant
* clusters the enrichment overwrites the tenant from the deployment identity anyway).
*/
protected static Optional<String> extractFleetTeamId(JsonNode afterField) {
if (afterField == null) {
return Optional.empty();
}
JsonNode teamId = afterField.get("team_id");
if (teamId == null || teamId.isNull()) {
return Optional.empty();
}
return Optional.of(teamId.asText());
}

/**
* Get effective timestamp for the event - uses event timestamp from source data if available,
* falls back to Debezium processing timestamp
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,10 +23,13 @@ public class DebeziumCassandraMessageHandler
extends DebeziumMessageHandler<UnifiedLogEvent, DeserializedDebeziumMessage> {

private final UnifiedLogEventRepository repository;
private final DebeziumEventValidator eventValidator;

public DebeziumCassandraMessageHandler(UnifiedLogEventRepository repository, ObjectMapper objectMapper) {
public DebeziumCassandraMessageHandler(UnifiedLogEventRepository repository, ObjectMapper objectMapper,
DebeziumEventValidator eventValidator) {
super(objectMapper);
this.repository = repository;
this.eventValidator = eventValidator;
}

@Override
Expand All @@ -39,6 +42,16 @@ public Destination getDestination() {
return Destination.CASSANDRA_EVENT_LOG;
}

/**
* Same tenant guard as the Kafka/Pinot handler: an event whose tenant could not be resolved
* (shared cluster — e.g. a Fleet CDC row without a stamped {@code team_id}, or a team with no
* tenant mapping) is dropped instead of being written.
*/
@Override
protected boolean isValidMessage(DeserializedDebeziumMessage message) {
return eventValidator.isValid(message);
}

Comment thread
coderabbitai[bot] marked this conversation as resolved.
@Override
protected UnifiedLogEvent transform(DeserializedDebeziumMessage debeziumMessage, IntegratedToolEnrichedData enrichedData) {
UnifiedLogEvent logEvent = new UnifiedLogEvent();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,10 @@ public class Activity {

@JsonProperty("user_email")
private String userEmail;


@JsonProperty("team_id")
private Long teamId;

// Enrichment fields (not from Debezium)
private String agentId; // From Redis cache or Fleet DB lookup by hostId
private Integer hostId; // From HostActivity join
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
package com.openframe.stream.service;

import com.openframe.data.model.enums.IntegratedToolType;

/**
* Resolves a canonical tenantId from a cluster-level identifier carried by an
* integrated-tool event (e.g. the MeshCentral {@code domain} in shared SaaS
Expand All @@ -16,4 +18,14 @@ public interface ClusterTenantIdResolver {
* @return canonical tenantId, or {@code null} when no tenant matches
*/
String resolveTenantId(String clusterName);

/**
* Tool-aware resolution: {@code key} is the tool-specific tenant discriminator
* (MeshCentral {@code domain}, Fleet {@code team_id}, …).
*
* @return canonical tenantId, or {@code null} when no tenant matches
*/
default String resolveTenantId(IntegratedToolType toolType, String key) {
return resolveTenantId(key);
}
}
Loading
Loading