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();

Reply via email to