diff --git a/dynamic-control/src/main/java/io/opentelemetry/contrib/dynamic/policy/AbstractPolicyProvider.java b/dynamic-control/src/main/java/io/opentelemetry/contrib/dynamic/policy/AbstractPolicyProvider.java new file mode 100644 index 000000000..1fee2d5d7 --- /dev/null +++ b/dynamic-control/src/main/java/io/opentelemetry/contrib/dynamic/policy/AbstractPolicyProvider.java @@ -0,0 +1,39 @@ +/* + * Copyright The OpenTelemetry Authors + * SPDX-License-Identifier: Apache-2.0 + */ + +package io.opentelemetry.contrib.dynamic.policy; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.Objects; +import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Consumer; + +/** Base {@link PolicyProvider} implementation for providers that maintain a current snapshot. */ +abstract class AbstractPolicyProvider implements PolicyProvider { + private final AtomicReference> currentPolicies = + new AtomicReference<>(Collections.emptyList()); + + protected final List getCurrentPolicies() { + return Objects.requireNonNull(currentPolicies.get(), "currentPolicies cannot be null"); + } + + protected final List updateCurrentPolicies(List policies) { + List snapshot = + Collections.unmodifiableList( + new ArrayList<>(Objects.requireNonNull(policies, "policies cannot be null"))); + currentPolicies.set(snapshot); + return snapshot; + } + + protected final List updateCurrentPoliciesAndNotify( + List policies, Consumer> onUpdate) { + Objects.requireNonNull(onUpdate, "onUpdate cannot be null"); + List snapshot = updateCurrentPolicies(policies); + onUpdate.accept(snapshot); + return snapshot; + } +} diff --git a/dynamic-control/src/main/java/io/opentelemetry/contrib/dynamic/policy/MappedPolicySourceConverter.java b/dynamic-control/src/main/java/io/opentelemetry/contrib/dynamic/policy/MappedPolicySourceConverter.java new file mode 100644 index 000000000..8c7000c8a --- /dev/null +++ b/dynamic-control/src/main/java/io/opentelemetry/contrib/dynamic/policy/MappedPolicySourceConverter.java @@ -0,0 +1,125 @@ +/* + * Copyright The OpenTelemetry Authors + * SPDX-License-Identifier: Apache-2.0 + */ + +package io.opentelemetry.contrib.dynamic.policy; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.fasterxml.jackson.databind.node.ObjectNode; +import io.opentelemetry.contrib.dynamic.policy.registry.PolicySourceMappingConfig; +import io.opentelemetry.contrib.dynamic.policy.source.JsonSourceWrapper; +import io.opentelemetry.contrib.dynamic.policy.source.KeyValueSourceWrapper; +import io.opentelemetry.contrib.dynamic.policy.source.SourceKind; +import io.opentelemetry.contrib.dynamic.policy.source.SourceWrapper; +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.Set; +import javax.annotation.Nullable; + +/** + * Converts parsed source policy entries into validated telemetry policies using source mappings. + */ +final class MappedPolicySourceConverter { + private static final ObjectMapper MAPPER = new ObjectMapper(); + + private final Map policyIdToPolicyType; + private final List validators; + + private MappedPolicySourceConverter( + Map policyIdToPolicyType, List validators) { + this.policyIdToPolicyType = + Collections.unmodifiableMap( + new HashMap<>( + Objects.requireNonNull( + policyIdToPolicyType, "policyIdToPolicyType cannot be null"))); + this.validators = + Collections.unmodifiableList( + new ArrayList<>(Objects.requireNonNull(validators, "validators cannot be null"))); + } + + static MappedPolicySourceConverter create( + List mappings, List validators) { + return new MappedPolicySourceConverter(buildPolicyIdToPolicyType(mappings), validators); + } + + Set getMappedPolicyIds() { + return policyIdToPolicyType.keySet(); + } + + List convert(List sources, SourceKind sourceKind) { + Objects.requireNonNull(sources, "sources cannot be null"); + Objects.requireNonNull(sourceKind, "sourceKind cannot be null"); + List policies = new ArrayList<>(); + for (SourceWrapper source : sources) { + TelemetryPolicy policy = convert(source, sourceKind); + if (policy != null) { + policies.add(policy); + } + } + return policies; + } + + @Nullable + TelemetryPolicy convert(SourceWrapper source, SourceKind sourceKind) { + Objects.requireNonNull(source, "source cannot be null"); + Objects.requireNonNull(sourceKind, "sourceKind cannot be null"); + String incomingPolicyId = source.getPolicyType(); + if (incomingPolicyId == null || incomingPolicyId.isEmpty()) { + return null; + } + String policyType = policyIdToPolicyType.get(incomingPolicyId); + if (policyType == null) { + return null; + } + SourceWrapper normalizedSource = remapSourcePolicyType(source, policyType); + if (normalizedSource == null) { + return null; + } + for (PolicyValidator validator : validators) { + if (!policyType.equals(validator.getPolicyType())) { + continue; + } + TelemetryPolicy policy = validator.validate(normalizedSource, sourceKind); + if (policy != null) { + return policy; + } + } + return null; + } + + private static Map buildPolicyIdToPolicyType( + List mappings) { + Objects.requireNonNull(mappings, "mappings cannot be null"); + Map mapping = new HashMap<>(); + for (PolicySourceMappingConfig item : mappings) { + mapping.put(item.getPolicyId(), item.getPolicyType()); + } + return mapping; + } + + @Nullable + private static SourceWrapper remapSourcePolicyType( + SourceWrapper source, String mappedPolicyType) { + if (source instanceof JsonSourceWrapper) { + JsonNode node = ((JsonSourceWrapper) source).asJsonNode(); + if (!node.isObject() || node.size() != 1) { + return null; + } + JsonNode value = node.elements().next(); + ObjectNode remappedNode = MAPPER.createObjectNode(); + remappedNode.set(mappedPolicyType, value); + return new JsonSourceWrapper(remappedNode); + } + if (source instanceof KeyValueSourceWrapper) { + KeyValueSourceWrapper keyValue = (KeyValueSourceWrapper) source; + return new KeyValueSourceWrapper(mappedPolicyType, keyValue.getValue()); + } + return source; + } +} diff --git a/dynamic-control/src/main/java/io/opentelemetry/contrib/dynamic/policy/OpampPolicyProvider.java b/dynamic-control/src/main/java/io/opentelemetry/contrib/dynamic/policy/OpampPolicyProvider.java index 7286135e2..ea6e68ec8 100644 --- a/dynamic-control/src/main/java/io/opentelemetry/contrib/dynamic/policy/OpampPolicyProvider.java +++ b/dynamic-control/src/main/java/io/opentelemetry/contrib/dynamic/policy/OpampPolicyProvider.java @@ -10,7 +10,6 @@ import com.fasterxml.jackson.databind.node.ObjectNode; import io.opentelemetry.contrib.dynamic.policy.registry.PolicySourceMappingConfig; import io.opentelemetry.contrib.dynamic.policy.source.JsonSourceWrapper; -import io.opentelemetry.contrib.dynamic.policy.source.KeyValueSourceWrapper; import io.opentelemetry.contrib.dynamic.policy.source.SourceFormat; import io.opentelemetry.contrib.dynamic.policy.source.SourceKind; import io.opentelemetry.contrib.dynamic.policy.source.SourceWrapper; @@ -53,7 +52,7 @@ * location key, maps incoming policy IDs to internal policy types, validates them with the supplied * validators, and publishes the resulting policies to callers. */ -public final class OpampPolicyProvider implements PolicyProvider { +public final class OpampPolicyProvider extends AbstractPolicyProvider { private static final Logger logger = Logger.getLogger(OpampPolicyProvider.class.getName()); private static final ObjectMapper MAPPER = new ObjectMapper(); @@ -75,11 +74,8 @@ public final class OpampPolicyProvider implements PolicyProvider { @Nullable private final String serviceEnvironment; private final Map headers; private final SourceFormat format; - private final List validators; - private final Map policyIdToPolicyType; + private final MappedPolicySourceConverter sourceConverter; private final MutablePeriodicDelay pollingDelay; - private final AtomicReference> currentPolicies = - new AtomicReference<>(Collections.emptyList()); private final AtomicReference clientRef = new AtomicReference<>(); private final AtomicReference shutdownHookRef = new AtomicReference<>(); @@ -114,22 +110,17 @@ public OpampPolicyProvider( this.serviceEnvironment = getServiceEnvironment(properties); this.headers = Collections.unmodifiableMap(new HashMap<>(properties.getMap(OPAMP_HEADERS))); this.format = format; - this.validators = Collections.unmodifiableList(new ArrayList<>(validators)); - this.policyIdToPolicyType = Collections.unmodifiableMap(buildPolicyIdToPolicyType(mappings)); + this.sourceConverter = MappedPolicySourceConverter.create(mappings, validators); this.pollingDelay = new MutablePeriodicDelay( Objects.requireNonNull( GLOBAL_POLLING_INTERVAL.get(), "polling interval cannot be null")); } - /** - * Returns the latest validated policies received from OpAMP. - * - *

The returned list is the current immutable snapshot held by this provider. - */ + /** Returns the latest validated policies received from OpAMP. */ @Override public List fetchPolicies() { - return Objects.requireNonNull(currentPolicies.get(), "currentPolicies cannot be null"); + return getCurrentPolicies(); } /** @@ -225,27 +216,21 @@ private RemoteConfigStatus handleMessage( } AgentConfigMap configMap = remoteConfig.config; if (configMap == null || configMap.config_map == null || configMap.config_map.isEmpty()) { - List empty = Collections.emptyList(); - currentPolicies.set(empty); - onUpdate.accept(empty); + updateCurrentPoliciesAndNotify(Collections.emptyList(), onUpdate); return buildStatus( RemoteConfigStatuses.RemoteConfigStatuses_FAILED, remoteConfig.config_hash); } AgentConfigFile selected = configMap.config_map.get(location); if (selected == null || selected.body == null) { logger.info("No OpAMP config payload found for location key: " + location); - List empty = Collections.emptyList(); - currentPolicies.set(empty); - onUpdate.accept(empty); + updateCurrentPoliciesAndNotify(Collections.emptyList(), onUpdate); return buildStatus( RemoteConfigStatuses.RemoteConfigStatuses_FAILED, remoteConfig.config_hash); } List policies = new ArrayList<>(); parsePolicyText(location, selected.body.utf8(), policies); - List snapshot = Collections.unmodifiableList(new ArrayList<>(policies)); - currentPolicies.set(snapshot); - onUpdate.accept(snapshot); + List snapshot = updateCurrentPoliciesAndNotify(policies, onUpdate); RemoteConfigStatuses status = snapshot.isEmpty() ? RemoteConfigStatuses.RemoteConfigStatuses_FAILED @@ -257,36 +242,13 @@ private void parsePolicyText(String key, String policyText, List parsedSources = format.parse(policyText); if (parsedSources == null && format == SourceFormat.JSONKEYVALUE) { - parsedSources = parseMappedJsonObject(policyText, policyIdToPolicyType.keySet()); + parsedSources = parseMappedJsonObject(policyText, sourceConverter.getMappedPolicyIds()); } if (parsedSources == null) { logger.info("Ignoring invalid OpAMP config entry for key: " + key); return; } - for (SourceWrapper source : parsedSources) { - String incomingPolicyId = source.getPolicyType(); - if (incomingPolicyId == null || incomingPolicyId.isEmpty()) { - continue; - } - String mappedPolicyType = policyIdToPolicyType.get(incomingPolicyId); - if (mappedPolicyType == null) { - continue; - } - SourceWrapper normalizedSource = remapSourcePolicyType(source, mappedPolicyType); - if (normalizedSource == null) { - continue; - } - for (PolicyValidator validator : validators) { - if (!mappedPolicyType.equals(validator.getPolicyType())) { - continue; - } - TelemetryPolicy policy = validator.validate(normalizedSource, SourceKind.OPAMP); - if (policy != null) { - out.add(policy); - break; - } - } - } + out.addAll(sourceConverter.convert(parsedSources, SourceKind.OPAMP)); } private void stop() { @@ -438,35 +400,6 @@ private static List parseMappedJsonObject( } } - private static Map buildPolicyIdToPolicyType( - List mappings) { - Map mapping = new HashMap<>(); - for (PolicySourceMappingConfig item : mappings) { - mapping.put(item.getPolicyId(), item.getPolicyType()); - } - return mapping; - } - - @Nullable - private static SourceWrapper remapSourcePolicyType( - SourceWrapper source, String mappedPolicyType) { - if (source instanceof JsonSourceWrapper) { - JsonNode node = ((JsonSourceWrapper) source).asJsonNode(); - if (!node.isObject() || node.size() != 1) { - return null; - } - JsonNode value = node.elements().next(); - ObjectNode remappedNode = MAPPER.createObjectNode(); - remappedNode.set(mappedPolicyType, value); - return new JsonSourceWrapper(remappedNode); - } - if (source instanceof KeyValueSourceWrapper) { - KeyValueSourceWrapper keyValue = (KeyValueSourceWrapper) source; - return new KeyValueSourceWrapper(mappedPolicyType, keyValue.getValue()); - } - return source; - } - private static RemoteConfigStatus buildStatus( RemoteConfigStatuses status, @Nullable ByteString hash) { if (hash != null && status == RemoteConfigStatuses.RemoteConfigStatuses_APPLIED) {