This is an automated email from the ASF dual-hosted git repository.
gyfora pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/flink-kubernetes-operator.git
The following commit(s) were added to refs/heads/main by this push:
new dd2faf9e [FLINK-40414] Autoscaler state ConfigMap is treated as
trusted input (#1184)
dd2faf9e is described below
commit dd2faf9e158f41551be82e70e8f5a9bf1c8e9282
Author: Dennis-Mircea Ciupitu <[email protected]>
AuthorDate: Wed Aug 19 11:47:55 2026 +0300
[FLINK-40414] Autoscaler state ConfigMap is treated as trusted input (#1184)
---
.../apache/flink/autoscaler/JobAutoScalerImpl.java | 33 +++++
.../flink/autoscaler/tuning/MemoryTuning.java | 34 +++++
.../flink/autoscaler/JobAutoScalerImplTest.java | 22 +++
.../flink/autoscaler/tuning/MemoryTuningTest.java | 4 +
.../state/KubernetesAutoScalerStateStore.java | 165 ++++++++++++---------
.../state/KubernetesAutoScalerStateStoreTest.java | 26 ++++
6 files changed, 215 insertions(+), 69 deletions(-)
diff --git
a/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/JobAutoScalerImpl.java
b/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/JobAutoScalerImpl.java
index d516a82a..c6154f1c 100644
---
a/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/JobAutoScalerImpl.java
+++
b/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/JobAutoScalerImpl.java
@@ -26,6 +26,7 @@ import
org.apache.flink.autoscaler.metrics.AutoscalerFlinkMetrics;
import org.apache.flink.autoscaler.realizer.ScalingRealizer;
import org.apache.flink.autoscaler.state.AutoScalerStateStore;
import org.apache.flink.autoscaler.tuning.ConfigChanges;
+import org.apache.flink.autoscaler.tuning.MemoryTuning;
import org.apache.flink.configuration.PipelineOptions;
import org.apache.flink.util.Preconditions;
@@ -35,7 +36,9 @@ import org.slf4j.LoggerFactory;
import java.time.Clock;
import java.util.HashMap;
import java.util.Map;
+import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
+import java.util.stream.Collectors;
import static
org.apache.flink.autoscaler.config.AutoScalerOptions.AUTOSCALER_ENABLED;
@@ -165,10 +168,40 @@ public class JobAutoScalerImpl<KEY, Context extends
JobAutoScalerContext<KEY>>
}
ConfigChanges configChanges = stateStore.getConfigChanges(ctx);
+ dropNonTunableKeys(ctx, configChanges);
LOG.debug("Applying config overrides: {}", configChanges);
scalingRealizer.realizeConfigOverrides(ctx, configChanges);
}
+ /**
+ * Drops any override or removal key that is not one memory tuning is
allowed to set. The config
+ * overrides are read back from the autoscaler state, which is writable by
any workload in the
+ * namespace, so without this filter a poisoned entry could inject
arbitrary Flink configuration
+ * (for example {@code env.java.opts} or a pod template) into the managed
deployment spec.
+ */
+ private void dropNonTunableKeys(Context ctx, ConfigChanges configChanges) {
+ Set<String> allowed = MemoryTuning.TUNABLE_CONFIG_KEYS;
+ Set<String> droppedOverrides =
+ configChanges.getOverrides().keySet().stream()
+ .filter(key -> !allowed.contains(key))
+ .collect(Collectors.toSet());
+ Set<String> droppedRemovals =
+ configChanges.getRemovals().stream()
+ .filter(key -> !allowed.contains(key))
+ .collect(Collectors.toSet());
+ if (droppedOverrides.isEmpty() && droppedRemovals.isEmpty()) {
+ return;
+ }
+ LOG.warn(
+ "Ignoring unexpected autoscaler config-override keys for {}:
overrides={}, removals={}. "
+ + "Only memory-tuning keys are applied.",
+ ctx.getJobKey(),
+ droppedOverrides,
+ droppedRemovals);
+ configChanges.getOverrides().keySet().removeAll(droppedOverrides);
+ configChanges.getRemovals().removeAll(droppedRemovals);
+ }
+
private void runScalingLogic(Context ctx, AutoscalerFlinkMetrics
autoscalerMetrics)
throws Exception {
var cycleState = ctx.getScalingCycleState();
diff --git
a/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/tuning/MemoryTuning.java
b/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/tuning/MemoryTuning.java
index 57ea353d..92eef873 100644
---
a/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/tuning/MemoryTuning.java
+++
b/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/tuning/MemoryTuning.java
@@ -29,6 +29,7 @@ import org.apache.flink.autoscaler.topology.JobTopology;
import org.apache.flink.autoscaler.topology.ShipStrategy;
import org.apache.flink.autoscaler.topology.VertexInfo;
import org.apache.flink.autoscaler.utils.ResourceCheckUtils;
+import org.apache.flink.configuration.ConfigOption;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.configuration.IllegalConfigurationException;
import org.apache.flink.configuration.MemorySize;
@@ -49,7 +50,9 @@ import org.slf4j.LoggerFactory;
import java.math.BigDecimal;
import java.math.RoundingMode;
import java.util.Arrays;
+import java.util.HashSet;
import java.util.Map;
+import java.util.Set;
import static
org.apache.flink.autoscaler.metrics.ScalingMetric.HEAP_MEMORY_USED;
import static
org.apache.flink.autoscaler.metrics.ScalingMetric.MANAGED_MEMORY_USED;
@@ -66,6 +69,37 @@ public class MemoryTuning {
private static final ConfigChanges EMPTY_CONFIG = new ConfigChanges();
+ /**
+ * The Flink configuration keys memory tuning is allowed to override or
remove. Used to filter
+ * the {@link ConfigChanges} read back from the autoscaler state, which is
writable by any
+ * workload in the namespace, so that a poisoned entry cannot inject
arbitrary Flink
+ * configuration (for example {@code env.java.opts} or a pod template)
into a managed
+ * deployment.
+ */
+ public static final Set<String> TUNABLE_CONFIG_KEYS =
collectTunableConfigKeys();
+
+ private static Set<String> collectTunableConfigKeys() {
+ // Must stay in sync with the keys tuneTaskManagerMemory emits
+ ConfigOption<?>[] options = {
+ TaskManagerOptions.TOTAL_PROCESS_MEMORY,
+ TaskManagerOptions.TOTAL_FLINK_MEMORY,
+ TaskManagerOptions.TASK_HEAP_MEMORY,
+ TaskManagerOptions.FRAMEWORK_HEAP_MEMORY,
+ TaskManagerOptions.MANAGED_MEMORY_FRACTION,
+ TaskManagerOptions.MANAGED_MEMORY_SIZE,
+ TaskManagerOptions.NETWORK_MEMORY_MIN,
+ TaskManagerOptions.NETWORK_MEMORY_MAX,
+ TaskManagerOptions.JVM_OVERHEAD_FRACTION,
+ TaskManagerOptions.JVM_METASPACE
+ };
+ Set<String> keys = new HashSet<>();
+ for (ConfigOption<?> option : options) {
+ keys.add(option.key());
+ option.fallbackKeys().forEach(fallbackKey ->
keys.add(fallbackKey.getKey()));
+ }
+ return Set.copyOf(keys);
+ }
+
/**
* Emits a Configuration which contains overrides for the current
configuration. We are not
* modifying the config directly, but we are emitting ConfigChanges which
contain any overrides
diff --git
a/flink-autoscaler/src/test/java/org/apache/flink/autoscaler/JobAutoScalerImplTest.java
b/flink-autoscaler/src/test/java/org/apache/flink/autoscaler/JobAutoScalerImplTest.java
index 7eac783a..43240b6a 100644
---
a/flink-autoscaler/src/test/java/org/apache/flink/autoscaler/JobAutoScalerImplTest.java
+++
b/flink-autoscaler/src/test/java/org/apache/flink/autoscaler/JobAutoScalerImplTest.java
@@ -416,6 +416,28 @@ public class JobAutoScalerImplTest {
assertTrue(sawConfig, "Config overrides must be re-applied while job
is not running");
}
+ @Test
+ void testConfigOverridesDropNonTunableKeys() throws Exception {
+
context.getConfiguration().set(AutoScalerOptions.MEMORY_TUNING_ENABLED, true);
+ var autoscaler =
+ new JobAutoScalerImpl<>(
+ null, null, null, eventCollector, scalingRealizer,
stateStore);
+
+ ConfigChanges poisoned = new ConfigChanges();
+ poisoned.addOverride(TaskManagerOptions.TOTAL_PROCESS_MEMORY.key(), "2
gb");
+ poisoned.addOverride("env.java.opts.all", "-XX:+SomethingEvil");
+ poisoned.getRemovals().add("kubernetes.pod-template-file");
+ stateStore.storeConfigChanges(context, poisoned);
+ stateStore.flush(context);
+
+ autoscaler.applyConfigOverrides(context);
+
+ var event = getEvent();
+ assertThat(event.getConfigChanges().getOverrides())
+
.containsOnlyKeys(TaskManagerOptions.TOTAL_PROCESS_MEMORY.key());
+ assertThat(event.getConfigChanges().getRemovals()).isEmpty();
+ }
+
@Test
void testApplyConfigOverrides() throws Exception {
context.getConfiguration().set(AutoScalerOptions.MEMORY_TUNING_ENABLED, true);
diff --git
a/flink-autoscaler/src/test/java/org/apache/flink/autoscaler/tuning/MemoryTuningTest.java
b/flink-autoscaler/src/test/java/org/apache/flink/autoscaler/tuning/MemoryTuningTest.java
index 35ccb633..4b95d6e7 100644
---
a/flink-autoscaler/src/test/java/org/apache/flink/autoscaler/tuning/MemoryTuningTest.java
+++
b/flink-autoscaler/src/test/java/org/apache/flink/autoscaler/tuning/MemoryTuningTest.java
@@ -133,6 +133,10 @@ public class MemoryTuningTest {
.next()
.getKey());
+ assertThat(configChanges.getOverrides().keySet())
+ .isSubsetOf(MemoryTuning.TUNABLE_CONFIG_KEYS);
+
assertThat(configChanges.getRemovals()).isSubsetOf(MemoryTuning.TUNABLE_CONFIG_KEYS);
+
assertThat(eventHandler.events.poll().getMessage())
.startsWith(
"Memory tuning recommends the following configuration
(automatic tuning is enabled):");
diff --git
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/autoscaler/state/KubernetesAutoScalerStateStore.java
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/autoscaler/state/KubernetesAutoScalerStateStore.java
index 2e989048..16c01e0e 100644
---
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/autoscaler/state/KubernetesAutoScalerStateStore.java
+++
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/autoscaler/state/KubernetesAutoScalerStateStore.java
@@ -57,6 +57,7 @@ import java.util.Map;
import java.util.Optional;
import java.util.SortedMap;
import java.util.TreeMap;
+import java.util.function.Supplier;
import java.util.zip.GZIPInputStream;
import java.util.zip.GZIPOutputStream;
@@ -116,38 +117,20 @@ public class KubernetesAutoScalerStateStore
@Override
public Map<JobVertexID, SortedMap<Instant, ScalingSummary>>
getScalingHistory(
KubernetesJobAutoScalerContext jobContext) {
- Optional<String> serializedScalingHistory =
- configMapStore.getSerializedState(jobContext,
SCALING_HISTORY_KEY);
- if (serializedScalingHistory.isEmpty()) {
- return new HashMap<>();
- }
- try {
- return deserializeScalingHistory(serializedScalingHistory.get());
- } catch (JacksonException e) {
- LOG.error(
- "Could not deserialize scaling history, possibly the
format changed. Discarding...",
- e);
- configMapStore.removeSerializedState(jobContext,
SCALING_HISTORY_KEY);
- return new HashMap<>();
- }
+ return readState(
+ jobContext,
+ SCALING_HISTORY_KEY,
+ KubernetesAutoScalerStateStore::deserializeScalingHistory,
+ HashMap::new);
}
@Override
public ScalingTracking getScalingTracking(KubernetesJobAutoScalerContext
jobContext) {
- Optional<String> serializedRescalingHistory =
- configMapStore.getSerializedState(jobContext,
SCALING_TRACKING_KEY);
- if (serializedRescalingHistory.isEmpty()) {
- return new ScalingTracking();
- }
- try {
- return
deserializeScalingTracking(serializedRescalingHistory.get());
- } catch (JacksonException e) {
- LOG.error(
- "Could not deseri alize rescaling history, possibly the
format changed. Discarding...",
- e);
- configMapStore.removeSerializedState(jobContext,
SCALING_TRACKING_KEY);
- return new ScalingTracking();
- }
+ return readState(
+ jobContext,
+ SCALING_TRACKING_KEY,
+ KubernetesAutoScalerStateStore::deserializeScalingTracking,
+ ScalingTracking::new);
}
@Override
@@ -167,20 +150,11 @@ public class KubernetesAutoScalerStateStore
@Override
public SortedMap<Instant, CollectedMetrics> getCollectedMetrics(
KubernetesJobAutoScalerContext jobContext) {
- Optional<String> serializedEvaluatedMetricsOpt =
- configMapStore.getSerializedState(jobContext,
COLLECTED_METRICS_KEY);
- if (serializedEvaluatedMetricsOpt.isEmpty()) {
- return new TreeMap<>();
- }
- try {
- return
deserializeEvaluatedMetrics(serializedEvaluatedMetricsOpt.get());
- } catch (JacksonException e) {
- LOG.error(
- "Could not deserialize metric history, possibly the format
changed. Discarding...",
- e);
- configMapStore.removeSerializedState(jobContext,
COLLECTED_METRICS_KEY);
- return new TreeMap<>();
- }
+ return readState(
+ jobContext,
+ COLLECTED_METRICS_KEY,
+ KubernetesAutoScalerStateStore::deserializeEvaluatedMetrics,
+ TreeMap::new);
}
@Override
@@ -200,19 +174,21 @@ public class KubernetesAutoScalerStateStore
@Nonnull
@Override
public Map<String, String>
getParallelismOverrides(KubernetesJobAutoScalerContext jobContext) {
- return configMapStore
- .getSerializedState(jobContext, PARALLELISM_OVERRIDES_KEY)
-
.map(KubernetesAutoScalerStateStore::deserializeParallelismOverrides)
- .orElse(new HashMap<>());
+ return readState(
+ jobContext,
+ PARALLELISM_OVERRIDES_KEY,
+ KubernetesAutoScalerStateStore::sanitizeParallelismOverrides,
+ HashMap::new);
}
@Nonnull
@Override
public ConfigChanges getConfigChanges(KubernetesJobAutoScalerContext
jobContext) {
- return configMapStore
- .getSerializedState(jobContext, CONFIG_OVERRIDES_KEY)
-
.map(KubernetesAutoScalerStateStore::deserializeConfigOverrides)
- .orElse(new ConfigChanges());
+ return readState(
+ jobContext,
+ CONFIG_OVERRIDES_KEY,
+ KubernetesAutoScalerStateStore::deserializeConfigOverrides,
+ ConfigChanges::new);
}
@Override
@@ -243,20 +219,46 @@ public class KubernetesAutoScalerStateStore
@Nonnull
@Override
public DelayedScaleDown getDelayedScaleDown(KubernetesJobAutoScalerContext
jobContext) {
- Optional<String> delayedScaleDown =
- configMapStore.getSerializedState(jobContext,
DELAYED_SCALE_DOWN);
- if (delayedScaleDown.isEmpty()) {
- return new DelayedScaleDown();
+ return readState(
+ jobContext,
+ DELAYED_SCALE_DOWN,
+ KubernetesAutoScalerStateStore::deserializeDelayedScaleDown,
+ DelayedScaleDown::new);
+ }
+
+ /** Deserializes a single stored state value, allowed to fail on untrusted
input. */
+ @FunctionalInterface
+ private interface StateDeserializer<T> {
+ T deserialize(String serialized) throws Exception;
+ }
+
+ /**
+ * Reads and deserializes one autoscaler state entry, treating the stored
value as untrusted.
+ * The autoscaler ConfigMap is writable by any workload in the namespace,
so any failure while
+ * reading an entry (unexpected schema, corrupt or oversized payload,
invalid value) discards
+ * that entry and falls back to empty, rather than propagating into the
reconcile loop or being
+ * acted upon. {@code Exception} is caught deliberately so no
deserialization path can escape
+ * this boundary.
+ */
+ private <T> T readState(
+ KubernetesJobAutoScalerContext jobContext,
+ String key,
+ StateDeserializer<T> deserializer,
+ Supplier<T> emptyValue) {
+ Optional<String> serialized =
configMapStore.getSerializedState(jobContext, key);
+ if (serialized.isEmpty()) {
+ return emptyValue.get();
}
-
try {
- return deserializeDelayedScaleDown(delayedScaleDown.get());
- } catch (JacksonException e) {
- LOG.warn(
- "Could not deserialize delayed scale down, possibly the
format changed. Discarding...",
+ return deserializer.deserialize(serialized.get());
+ } catch (Exception e) {
+ LOG.error(
+ "Discarding invalid autoscaler state '{}' for {}.",
+ key,
+ jobContext.getJobKey(),
e);
- configMapStore.removeSerializedState(jobContext,
DELAYED_SCALE_DOWN);
- return new DelayedScaleDown();
+ configMapStore.removeSerializedState(jobContext, key);
+ return emptyValue.get();
}
}
@@ -312,6 +314,36 @@ public class KubernetesAutoScalerStateStore
return ConfigurationUtils.convertValue(overrides, Map.class);
}
+ /**
+ * Deserializes the parallelism overrides and drops any entry whose value
is not a positive
+ * integer, since a vertex parallelism below one is never valid and the
stored value is
+ * untrusted. A malformed map as a whole still fails deserialization and
is discarded upstream.
+ */
+ private static Map<String, String> sanitizeParallelismOverrides(String
serialized) {
+ Map<String, String> overrides =
deserializeParallelismOverrides(serialized);
+ Map<String, String> sanitized = new HashMap<>();
+ overrides.forEach(
+ (vertexId, parallelism) -> {
+ if (isPositiveInt(parallelism)) {
+ sanitized.put(vertexId, parallelism);
+ } else {
+ LOG.warn(
+ "Dropping invalid parallelism override {}={}
from autoscaler state.",
+ vertexId,
+ parallelism);
+ }
+ });
+ return sanitized;
+ }
+
+ private static boolean isPositiveInt(String value) {
+ try {
+ return Integer.parseInt(value.trim()) > 0;
+ } catch (NumberFormatException e) {
+ return false;
+ }
+ }
+
@Nullable
private static String serializeConfigOverrides(ConfigChanges
configChanges) {
try {
@@ -322,14 +354,9 @@ public class KubernetesAutoScalerStateStore
}
}
- @Nullable
- private static ConfigChanges deserializeConfigOverrides(String
configOverrides) {
- try {
- return YAML_MAPPER.readValue(configOverrides, new
TypeReference<>() {});
- } catch (Exception e) {
- LOG.error("Failed to deserialize ConfigOverrides", e);
- return null;
- }
+ private static ConfigChanges deserializeConfigOverrides(String
configOverrides)
+ throws JacksonException {
+ return YAML_MAPPER.readValue(configOverrides, new TypeReference<>()
{});
}
private static String serializeDelayedScaleDown(DelayedScaleDown
delayedScaleDown)
diff --git
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/autoscaler/state/KubernetesAutoScalerStateStoreTest.java
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/autoscaler/state/KubernetesAutoScalerStateStoreTest.java
index 468b91f7..3096ad07 100644
---
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/autoscaler/state/KubernetesAutoScalerStateStoreTest.java
+++
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/autoscaler/state/KubernetesAutoScalerStateStoreTest.java
@@ -265,6 +265,32 @@ public class KubernetesAutoScalerStateStoreTest
.isEmpty();
}
+ @Test
+ void testDiscardMalformedConfigOverrides() throws Exception {
+ configMapStore.putSerializedState(
+ ctx, KubernetesAutoScalerStateStore.CONFIG_OVERRIDES_KEY,
"not-config-changes");
+ assertThat(stateStore.getConfigChanges(ctx).getOverrides()).isEmpty();
+ assertThat(
+ configMapStore.getSerializedState(
+ ctx,
KubernetesAutoScalerStateStore.CONFIG_OVERRIDES_KEY))
+ .isEmpty();
+ }
+
+ @Test
+ void testParallelismOverridesFailClosedAndSanitized() throws Exception {
+ configMapStore.putSerializedState(
+ ctx, KubernetesAutoScalerStateStore.PARALLELISM_OVERRIDES_KEY,
"###");
+ assertThat(stateStore.getParallelismOverrides(ctx)).isEmpty();
+
+ var v1 = new JobVertexID().toString();
+ var v2 = new JobVertexID().toString();
+ var v3 = new JobVertexID().toString();
+ stateStore.storeParallelismOverrides(ctx, Map.of(v1, "4", v2, "-3",
v3, "x"));
+ stateStore.flush(ctx);
+ assertThat(stateStore.getParallelismOverrides(ctx))
+ .containsExactlyInAnyOrderEntriesOf(Map.of(v1, "4"));
+ }
+
@Test
protected void testDiscardAllState() throws Exception {
super.testDiscardAllState();