This is an automated email from the ASF dual-hosted git repository.

gyfora pushed a commit to branch release-1.16
in repository https://gitbox.apache.org/repos/asf/flink-kubernetes-operator.git

commit c7ad36ed3d3a93ff2e620a86febe187f7f6eab31
Author: Purushottam Sinha <[email protected]>
AuthorDate: Wed Aug 19 12:10:43 2026 +0530

    [FLINK-40401] Harden autoscaler state deserialization in kubernetes 
operator. Bound decompressed data size. (#1181)
    
    Generated-by: Claude Code
---
 .../state/KubernetesAutoScalerStateStore.java      | 25 ++++++++++++--
 .../state/KubernetesAutoScalerStateStoreTest.java  | 39 ++++++++++++++++++++++
 2 files changed, 62 insertions(+), 2 deletions(-)

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 ed6549ad..2e989048 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
@@ -39,7 +39,6 @@ import 
org.apache.flink.shaded.jackson2.org.yaml.snakeyaml.LoaderOptions;
 
 import io.javaoperatorsdk.operator.processing.event.ResourceID;
 import lombok.SneakyThrows;
-import org.apache.commons.io.IOUtils;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -49,6 +48,7 @@ import javax.annotation.Nullable;
 import java.io.ByteArrayInputStream;
 import java.io.ByteArrayOutputStream;
 import java.io.IOException;
+import java.io.InputStream;
 import java.nio.charset.StandardCharsets;
 import java.time.Instant;
 import java.util.Base64;
@@ -80,6 +80,11 @@ public class KubernetesAutoScalerStateStore
 
     @VisibleForTesting protected static final int MAX_CM_BYTES = 1000000;
 
+    /* Caps decompressed value size against gzip bombs. Legit values are a few 
MB at most
+     * (MAX_CM_BYTES compressed, ~2.5-4x ratio); this is well below the YAML 
loader's limit,
+     * which guards the parser, not decompression memory. */
+    @VisibleForTesting protected static final int MAX_DECOMPRESSED_BYTES = 8 * 
1024 * 1024;
+
     protected static final ObjectMapper YAML_MAPPER =
             new ObjectMapper(yamlFactory())
                     .registerModule(new JavaTimeModule())
@@ -388,7 +393,7 @@ public class KubernetesAutoScalerStateStore
         try {
             byte[] bytes = Base64.getDecoder().decode(compressed);
             try (var zi = new GZIPInputStream(new 
ByteArrayInputStream(bytes))) {
-                return IOUtils.toString(zi, StandardCharsets.UTF_8);
+                return readBounded(zi, MAX_DECOMPRESSED_BYTES);
             }
         } catch (Exception e) {
             LOG.warn("Error while decompressing scaling data, treating as 
uncompressed");
@@ -397,6 +402,22 @@ public class KubernetesAutoScalerStateStore
         }
     }
 
+    private static String readBounded(InputStream in, int maxBytes) throws 
IOException {
+        var out = new ByteArrayOutputStream();
+        byte[] buffer = new byte[8192];
+        int totalRead = 0;
+        int read;
+        while ((read = in.read(buffer)) != -1) {
+            totalRead += read;
+            if (totalRead > maxBytes) {
+                throw new IOException(
+                        "Refusing to decompress data larger than " + maxBytes 
+ " bytes");
+            }
+            out.write(buffer, 0, read);
+        }
+        return out.toString(StandardCharsets.UTF_8);
+    }
+
     private static YAMLFactory yamlFactory() {
         // Set yaml size limit to 10mb
         var loaderOptions = new LoaderOptions();
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 b704cfcf..468b91f7 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
@@ -32,13 +32,16 @@ import 
io.javaoperatorsdk.operator.processing.event.ResourceID;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
 
+import java.io.ByteArrayOutputStream;
 import java.time.Instant;
+import java.util.Base64;
 import java.util.Collections;
 import java.util.HashMap;
 import java.util.Map;
 import java.util.Random;
 import java.util.SortedMap;
 import java.util.TreeMap;
+import java.util.zip.GZIPOutputStream;
 
 import static 
org.apache.flink.autoscaler.metrics.ScalingHistoryUtils.addToScalingHistoryAndStore;
 import static 
org.apache.flink.autoscaler.metrics.ScalingHistoryUtils.getTrimmedScalingHistory;
@@ -226,6 +229,42 @@ public class KubernetesAutoScalerStateStoreTest
                 .isEmpty();
     }
 
+    @Test
+    void testDiscardOversizedCompressedHistory() throws Exception {
+        // Craft a gzip "bomb": a tiny compressed payload (all zero bytes 
compress extremely
+        // well) that decompresses to more than MAX_DECOMPRESSED_BYTES, 
simulating a
+        // maliciously crafted ConfigMap entry. The store must reject it via 
bounded reading
+        // instead of decompressing it fully into memory.
+        var bomb = new ByteArrayOutputStream();
+        try (var gzip = new GZIPOutputStream(bomb)) {
+            byte[] chunk = new byte[1024 * 1024];
+            int chunks = 
(KubernetesAutoScalerStateStore.MAX_DECOMPRESSED_BYTES / chunk.length) + 1;
+            for (int i = 0; i < chunks; i++) {
+                gzip.write(chunk);
+            }
+        }
+        String bombPayload = 
Base64.getEncoder().encodeToString(bomb.toByteArray());
+
+        configMapStore.putSerializedState(
+                ctx, KubernetesAutoScalerStateStore.COLLECTED_METRICS_KEY, 
bombPayload);
+        configMapStore.putSerializedState(
+                ctx, KubernetesAutoScalerStateStore.SCALING_HISTORY_KEY, 
bombPayload);
+
+        var now = Instant.now();
+
+        assertThat(stateStore.getCollectedMetrics(ctx)).isEmpty();
+        assertThat(
+                        configMapStore.getSerializedState(
+                                ctx, 
KubernetesAutoScalerStateStore.COLLECTED_METRICS_KEY))
+                .isEmpty();
+
+        Assertions.assertEquals(new TreeMap<>(), 
getTrimmedScalingHistory(stateStore, ctx, now));
+        assertThat(
+                        configMapStore.getSerializedState(
+                                ctx, 
KubernetesAutoScalerStateStore.SCALING_HISTORY_KEY))
+                .isEmpty();
+    }
+
     @Test
     protected void testDiscardAllState() throws Exception {
         super.testDiscardAllState();

Reply via email to