Skip to content
Open
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
@@ -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<List<TelemetryPolicy>> currentPolicies =
new AtomicReference<>(Collections.<TelemetryPolicy>emptyList());

protected final List<TelemetryPolicy> getCurrentPolicies() {
return Objects.requireNonNull(currentPolicies.get(), "currentPolicies cannot be null");
}

protected final List<TelemetryPolicy> updateCurrentPolicies(List<TelemetryPolicy> policies) {
List<TelemetryPolicy> snapshot =
Collections.unmodifiableList(
new ArrayList<>(Objects.requireNonNull(policies, "policies cannot be null")));
currentPolicies.set(snapshot);
return snapshot;
}

protected final List<TelemetryPolicy> updateCurrentPoliciesAndNotify(
List<TelemetryPolicy> policies, Consumer<List<TelemetryPolicy>> onUpdate) {
Objects.requireNonNull(onUpdate, "onUpdate cannot be null");
List<TelemetryPolicy> snapshot = updateCurrentPolicies(policies);
onUpdate.accept(snapshot);
return snapshot;
}
}
Original file line number Diff line number Diff line change
@@ -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<String, String> policyIdToPolicyType;
private final List<PolicyValidator> validators;

private MappedPolicySourceConverter(
Map<String, String> policyIdToPolicyType, List<PolicyValidator> 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<PolicySourceMappingConfig> mappings, List<PolicyValidator> validators) {
return new MappedPolicySourceConverter(buildPolicyIdToPolicyType(mappings), validators);
}

Set<String> getMappedPolicyIds() {
return policyIdToPolicyType.keySet();
}

List<TelemetryPolicy> convert(List<SourceWrapper> sources, SourceKind sourceKind) {
Objects.requireNonNull(sources, "sources cannot be null");
Objects.requireNonNull(sourceKind, "sourceKind cannot be null");
List<TelemetryPolicy> 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<String, String> buildPolicyIdToPolicyType(
List<PolicySourceMappingConfig> mappings) {
Objects.requireNonNull(mappings, "mappings cannot be null");
Map<String, String> 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;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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();

Expand All @@ -75,11 +74,8 @@ public final class OpampPolicyProvider implements PolicyProvider {
@Nullable private final String serviceEnvironment;
private final Map<String, String> headers;
private final SourceFormat format;
private final List<PolicyValidator> validators;
private final Map<String, String> policyIdToPolicyType;
private final MappedPolicySourceConverter sourceConverter;
private final MutablePeriodicDelay pollingDelay;
private final AtomicReference<List<TelemetryPolicy>> currentPolicies =
new AtomicReference<>(Collections.<TelemetryPolicy>emptyList());
private final AtomicReference<OpampClient> clientRef = new AtomicReference<>();
private final AtomicReference<Thread> shutdownHookRef = new AtomicReference<>();

Expand Down Expand Up @@ -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.
*
* <p>The returned list is the current immutable snapshot held by this provider.
*/
/** Returns the latest validated policies received from OpAMP. */
@Override
public List<TelemetryPolicy> fetchPolicies() {
return Objects.requireNonNull(currentPolicies.get(), "currentPolicies cannot be null");
return getCurrentPolicies();
}

/**
Expand Down Expand Up @@ -225,27 +216,21 @@ private RemoteConfigStatus handleMessage(
}
AgentConfigMap configMap = remoteConfig.config;
if (configMap == null || configMap.config_map == null || configMap.config_map.isEmpty()) {
List<TelemetryPolicy> empty = Collections.emptyList();
currentPolicies.set(empty);
onUpdate.accept(empty);
updateCurrentPoliciesAndNotify(Collections.<TelemetryPolicy>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<TelemetryPolicy> empty = Collections.emptyList();
currentPolicies.set(empty);
onUpdate.accept(empty);
updateCurrentPoliciesAndNotify(Collections.<TelemetryPolicy>emptyList(), onUpdate);
return buildStatus(
RemoteConfigStatuses.RemoteConfigStatuses_FAILED, remoteConfig.config_hash);
}

List<TelemetryPolicy> policies = new ArrayList<>();
parsePolicyText(location, selected.body.utf8(), policies);
List<TelemetryPolicy> snapshot = Collections.unmodifiableList(new ArrayList<>(policies));
currentPolicies.set(snapshot);
onUpdate.accept(snapshot);
List<TelemetryPolicy> snapshot = updateCurrentPoliciesAndNotify(policies, onUpdate);
RemoteConfigStatuses status =
snapshot.isEmpty()
? RemoteConfigStatuses.RemoteConfigStatuses_FAILED
Expand All @@ -257,36 +242,13 @@ private void parsePolicyText(String key, String policyText, List<TelemetryPolicy
logger.info("Received OpAMP policy payload for key '" + key + "': " + policyText);
List<SourceWrapper> 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() {
Expand Down Expand Up @@ -438,35 +400,6 @@ private static List<SourceWrapper> parseMappedJsonObject(
}
}

private static Map<String, String> buildPolicyIdToPolicyType(
List<PolicySourceMappingConfig> mappings) {
Map<String, String> 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) {
Expand Down
Loading