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

rmetzger 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 39de41c4 [FLINK-40400] Harden job artifact fetching in kubernetes 
operator (#1180)
39de41c4 is described below

commit 39de41c4195b23f5bb0dd6868fe4bbe01234c9e1
Author: Purushottam Sinha <[email protected]>
AuthorDate: Wed Aug 19 13:26:21 2026 +0530

    [FLINK-40400] Harden job artifact fetching in kubernetes operator (#1180)
    
    * [FLINK-40400] Harden job artifact fetching in kubernetes operator. Ensure 
proper artifact validation.
    
    
    
    Generated-by: Claude Code
---
 .../operator/artifact/ArtifactManager.java         |  13 +-
 .../operator/artifact/HttpArtifactFetcher.java     | 128 +++++++++++++-
 .../operator/utils/JarUriValidationUtils.java      | 102 ++++++++++++
 .../operator/validation/DefaultValidator.java      |  59 +------
 .../operator/artifact/ArtifactManagerTest.java     | 184 ++++++++++++++++++++-
 5 files changed, 420 insertions(+), 66 deletions(-)

diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/artifact/ArtifactManager.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/artifact/ArtifactManager.java
index f40ae67b..7e229573 100644
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/artifact/ArtifactManager.java
+++ 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/artifact/ArtifactManager.java
@@ -20,6 +20,7 @@ package org.apache.flink.kubernetes.operator.artifact;
 import org.apache.flink.configuration.Configuration;
 import org.apache.flink.kubernetes.operator.api.spec.FlinkSessionJobSpec;
 import org.apache.flink.kubernetes.operator.config.FlinkConfigManager;
+import 
org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions;
 import org.apache.flink.util.FlinkRuntimeException;
 
 import io.fabric8.kubernetes.api.model.ObjectMeta;
@@ -58,7 +59,17 @@ public class ArtifactManager {
         createIfNotExists(targetDir);
         URI uri = new URI(jarURI);
         if ("http".equals(uri.getScheme()) || "https".equals(uri.getScheme())) 
{
-            return HttpArtifactFetcher.INSTANCE.fetch(jarURI, 
flinkConfiguration, targetDir);
+            // Take the scheme/host policy from the operator config (matching 
DefaultValidator);
+            // clone so the caller's config is not mutated.
+            var operatorConfig = configManager.getOperatorConfiguration();
+            var fetchConfig = flinkConfiguration.clone();
+            fetchConfig.set(
+                    KubernetesOperatorConfigOptions.JAR_URI_ALLOWED_SCHEMES,
+                    operatorConfig.getJarUriAllowedSchemes());
+            fetchConfig.set(
+                    
KubernetesOperatorConfigOptions.JAR_URI_DISALLOW_RESTRICTED_HOSTS,
+                    operatorConfig.isJarUriDisallowRestrictedHosts());
+            return HttpArtifactFetcher.INSTANCE.fetch(jarURI, fetchConfig, 
targetDir);
         } else {
             return FileSystemBasedArtifactFetcher.INSTANCE.fetch(
                     jarURI, flinkConfiguration, targetDir);
diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/artifact/HttpArtifactFetcher.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/artifact/HttpArtifactFetcher.java
index c2c6c2db..09e517c6 100644
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/artifact/HttpArtifactFetcher.java
+++ 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/artifact/HttpArtifactFetcher.java
@@ -19,6 +19,7 @@ package org.apache.flink.kubernetes.operator.artifact;
 
 import org.apache.flink.configuration.Configuration;
 import 
org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions;
+import org.apache.flink.kubernetes.operator.utils.JarUriValidationUtils;
 
 import org.apache.commons.io.FileUtils;
 import org.apache.commons.io.FilenameUtils;
@@ -26,38 +27,141 @@ import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 import java.io.File;
+import java.io.IOException;
 import java.net.HttpURLConnection;
+import java.net.MalformedURLException;
 import java.net.URL;
 import java.util.Map;
 
-/** Download the jar from the http resource. */
+/**
+ * Download the jar from the http resource. The scheme allowlist and 
restricted-host policy are read
+ * from the given configuration; {@link ArtifactManager} sets them from the 
operator configuration
+ * before calling.
+ */
 public class HttpArtifactFetcher implements ArtifactFetcher {
 
     public static final Logger LOG = 
LoggerFactory.getLogger(HttpArtifactFetcher.class);
     public static final HttpArtifactFetcher INSTANCE = new 
HttpArtifactFetcher();
 
+    // Maximum number of redirects to follow before giving up.
+    private static final int MAX_REDIRECTS = 5;
+
     @Override
     public File fetch(String uri, Configuration flinkConfiguration, File 
targetDir)
             throws Exception {
         var start = System.currentTimeMillis();
-        URL url = new URL(uri);
-        HttpURLConnection conn = (HttpURLConnection) url.openConnection();
+
+        // Scheme allowlist and restricted-host policy, set by ArtifactManager 
from the operator
+        // configuration.
+        var allowedSchemes =
+                
flinkConfiguration.get(KubernetesOperatorConfigOptions.JAR_URI_ALLOWED_SCHEMES);
+        var disallowRestrictedHosts =
+                flinkConfiguration.get(
+                        
KubernetesOperatorConfigOptions.JAR_URI_DISALLOW_RESTRICTED_HOSTS);
 
         // merged session job level header and cluster level header, session 
job level header take
         // precedence.
         Map<String, String> headers =
                 
flinkConfiguration.get(KubernetesOperatorConfigOptions.JAR_ARTIFACT_HTTP_HEADER);
 
-        if (headers != null) {
-            headers.forEach(conn::setRequestProperty);
-        }
+        // Follow redirects manually so each hop is validated against the same 
policy as the
+        // original URI.
+        String currentUri = uri;
+        URL originalUrl = null;
+        URL currentUrl;
+        HttpURLConnection conn;
+        int redirects = 0;
+        while (true) {
+            var validationError =
+                    JarUriValidationUtils.validateJarURI(
+                            currentUri, allowedSchemes, 
disallowRestrictedHosts);
+            if (validationError.isPresent()) {
+                throw new IOException(
+                        "Refusing to fetch artifact from '"
+                                + currentUri
+                                + "': "
+                                + validationError.get());
+            }
+
+            currentUrl = new URL(currentUri);
+            if (originalUrl == null) {
+                originalUrl = currentUrl;
+            }
+            conn = (HttpURLConnection) currentUrl.openConnection();
+            conn.setInstanceFollowRedirects(false);
+            // Only send the configured headers to the original host; drop 
them on a cross-host
+            // redirect.
+            if (headers != null && 
originalUrl.getHost().equalsIgnoreCase(currentUrl.getHost())) {
+                headers.forEach(conn::setRequestProperty);
+            }
+            conn.setRequestMethod("GET");
 
-        conn.setRequestMethod("GET");
+            // Release the connection on every path except the final 
(non-redirect) one, whose body
+            // is streamed below. This covers getResponseCode() and the 
redirect handling throwing.
+            boolean keepConnection = false;
+            try {
+                int status = conn.getResponseCode();
+                if (!isRedirect(status)) {
+                    keepConnection = true;
+                    break;
+                }
 
-        String fileName = FilenameUtils.getName(url.getPath());
+                String location = conn.getHeaderField("Location");
+                if (location == null || location.isEmpty()) {
+                    throw new IOException(
+                            "Received redirect (status "
+                                    + status
+                                    + ") from '"
+                                    + currentUri
+                                    + "' without a Location header");
+                }
+                if (++redirects > MAX_REDIRECTS) {
+                    throw new IOException(
+                            "Too many redirects (>"
+                                    + MAX_REDIRECTS
+                                    + ") while fetching artifact from '"
+                                    + uri
+                                    + "'");
+                }
+                URL nextUrl;
+                try {
+                    nextUrl = new URL(currentUrl, location);
+                } catch (MalformedURLException e) {
+                    throw new IOException(
+                            "Refusing to follow redirect from '"
+                                    + currentUri
+                                    + "' to '"
+                                    + location
+                                    + "': "
+                                    + e.getMessage());
+                }
+                // An HTTP fetch only follows http(s) redirects, even if other 
schemes (e.g. s3,
+                // hdfs) are in the jarURI allowlist for top-level use.
+                var nextScheme = nextUrl.getProtocol();
+                if (!"http".equalsIgnoreCase(nextScheme) && 
!"https".equalsIgnoreCase(nextScheme)) {
+                    throw new IOException(
+                            "Refusing to follow redirect from '"
+                                    + currentUri
+                                    + "' to non-http(s) target '"
+                                    + nextUrl
+                                    + "'");
+                }
+                currentUri = nextUrl.toString();
+            } finally {
+                if (!keepConnection) {
+                    conn.disconnect();
+                }
+            }
+        }
+
+        // Name the file from the original jarURI, not the redirect target, so 
a redirect can't
+        // change it (e.g. drop the .jar extension the JobManager upload 
requires).
+        String fileName = FilenameUtils.getName(originalUrl.getPath());
         File targetFile = new File(targetDir, fileName);
         try (var inputStream = conn.getInputStream()) {
             FileUtils.copyToFile(inputStream, targetFile);
+        } finally {
+            conn.disconnect();
         }
         LOG.debug(
                 "Copied file from {} to {}, cost {} ms",
@@ -66,4 +170,12 @@ public class HttpArtifactFetcher implements ArtifactFetcher 
{
                 System.currentTimeMillis() - start);
         return targetFile;
     }
+
+    private static boolean isRedirect(int status) {
+        return status == HttpURLConnection.HTTP_MOVED_PERM
+                || status == HttpURLConnection.HTTP_MOVED_TEMP
+                || status == HttpURLConnection.HTTP_SEE_OTHER
+                || status == 307 // Temporary Redirect
+                || status == 308; // Permanent Redirect
+    }
 }
diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/JarUriValidationUtils.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/JarUriValidationUtils.java
new file mode 100644
index 00000000..45b2c0f1
--- /dev/null
+++ 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/JarUriValidationUtils.java
@@ -0,0 +1,102 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.flink.kubernetes.operator.utils;
+
+import 
org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions;
+
+import java.net.InetAddress;
+import java.net.URI;
+import java.net.URISyntaxException;
+import java.net.UnknownHostException;
+import java.util.Collection;
+import java.util.Locale;
+import java.util.Optional;
+import java.util.Set;
+import java.util.stream.Collectors;
+
+/**
+ * Shared jarURI validation (scheme allowlist plus restricted-host checks), 
used both at
+ * admission/reconcile time and to re-validate every hop an artifact fetch is 
redirected through.
+ */
+public final class JarUriValidationUtils {
+
+    private JarUriValidationUtils() {}
+
+    public static Optional<String> validateJarURI(
+            String jarURI, Collection<String> allowedSchemes, boolean 
disallowRestrictedHosts) {
+        if (jarURI == null) {
+            return Optional.empty();
+        }
+
+        URI uri;
+        try {
+            uri = new URI(jarURI);
+        } catch (URISyntaxException e) {
+            return Optional.of("jarURI is not a valid URI: " + e.getMessage());
+        }
+
+        String scheme = uri.getScheme();
+        if (scheme == null) {
+            return Optional.of("jarURI must include a scheme");
+        }
+
+        Set<String> normalizedAllowedSchemes =
+                allowedSchemes.stream()
+                        .map(s -> s.toLowerCase(Locale.ROOT))
+                        .collect(Collectors.toSet());
+        if 
(!normalizedAllowedSchemes.contains(scheme.toLowerCase(Locale.ROOT))) {
+            return Optional.of(
+                    String.format(
+                            "jarURI scheme '%s' is not in the allowlist %s. 
Configure '%s' to extend the allowlist.",
+                            scheme,
+                            normalizedAllowedSchemes,
+                            
KubernetesOperatorConfigOptions.JAR_URI_ALLOWED_SCHEMES.key()));
+        }
+
+        if (("http".equalsIgnoreCase(scheme) || 
"https".equalsIgnoreCase(scheme))
+                && disallowRestrictedHosts) {
+            String host = uri.getHost();
+            if (host == null || host.isEmpty()) {
+                return Optional.of("jarURI must include a host for http/https 
schemes");
+            }
+            InetAddress[] addresses;
+            try {
+                // Check every resolved address, not just the first, since a 
host can resolve to
+                // multiple A/AAAA records.
+                addresses = InetAddress.getAllByName(host);
+            } catch (UnknownHostException e) {
+                return Optional.of("jarURI host '" + host + "' cannot be 
resolved");
+            }
+            for (InetAddress addr : addresses) {
+                if (isRestricted(addr)) {
+                    return Optional.of(
+                            "jarURI host '" + host + "' resolves to a 
restricted address");
+                }
+            }
+        }
+        return Optional.empty();
+    }
+
+    private static boolean isRestricted(InetAddress addr) {
+        return addr.isLoopbackAddress()
+                || addr.isLinkLocalAddress()
+                || addr.isSiteLocalAddress()
+                || addr.isAnyLocalAddress()
+                || addr.isMulticastAddress();
+    }
+}
diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/validation/DefaultValidator.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/validation/DefaultValidator.java
index 20b3c5c8..7b7186b2 100644
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/validation/DefaultValidator.java
+++ 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/validation/DefaultValidator.java
@@ -48,6 +48,7 @@ import 
org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptio
 import org.apache.flink.kubernetes.operator.exception.ReconciliationException;
 import org.apache.flink.kubernetes.operator.utils.FlinkStateSnapshotUtils;
 import org.apache.flink.kubernetes.operator.utils.IngressUtils;
+import org.apache.flink.kubernetes.operator.utils.JarUriValidationUtils;
 import org.apache.flink.kubernetes.operator.utils.ResourceConfigUtils;
 import org.apache.flink.kubernetes.utils.Constants;
 import org.apache.flink.runtime.clusterframework.TaskExecutorProcessUtils;
@@ -62,20 +63,14 @@ import org.slf4j.LoggerFactory;
 
 import javax.annotation.Nullable;
 
-import java.net.InetAddress;
-import java.net.URI;
-import java.net.URISyntaxException;
-import java.net.UnknownHostException;
 import java.util.ArrayList;
 import java.util.Collection;
 import java.util.List;
-import java.util.Locale;
 import java.util.Map;
 import java.util.Optional;
 import java.util.Set;
 import java.util.regex.Matcher;
 import java.util.regex.Pattern;
-import java.util.stream.Collectors;
 
 /** Default validator implementation for {@link FlinkDeployment}. */
 public class DefaultValidator implements FlinkResourceValidator {
@@ -316,56 +311,8 @@ public class DefaultValidator implements 
FlinkResourceValidator {
     @VisibleForTesting
     static Optional<String> validateJarURI(
             String jarURI, Collection<String> allowedSchemes, boolean 
disallowRestrictedHosts) {
-        if (jarURI == null) {
-            return Optional.empty();
-        }
-
-        URI uri;
-        try {
-            uri = new URI(jarURI);
-        } catch (URISyntaxException e) {
-            return Optional.of("jarURI is not a valid URI: " + e.getMessage());
-        }
-
-        String scheme = uri.getScheme();
-        if (scheme == null) {
-            return Optional.of("jarURI must include a scheme");
-        }
-
-        Set<String> normalizedAllowedSchemes =
-                allowedSchemes.stream()
-                        .map(s -> s.toLowerCase(Locale.ROOT))
-                        .collect(Collectors.toSet());
-        if 
(!normalizedAllowedSchemes.contains(scheme.toLowerCase(Locale.ROOT))) {
-            return Optional.of(
-                    String.format(
-                            "jarURI scheme '%s' is not in the allowlist %s. 
Configure '%s' to extend the allowlist.",
-                            scheme,
-                            normalizedAllowedSchemes,
-                            
KubernetesOperatorConfigOptions.JAR_URI_ALLOWED_SCHEMES.key()));
-        }
-
-        if (("http".equalsIgnoreCase(scheme) || 
"https".equalsIgnoreCase(scheme))
-                && disallowRestrictedHosts) {
-            String host = uri.getHost();
-            if (host == null || host.isEmpty()) {
-                return Optional.of("jarURI must include a host for http/https 
schemes");
-            }
-            InetAddress addr;
-            try {
-                addr = InetAddress.getByName(host);
-            } catch (UnknownHostException e) {
-                return Optional.of("jarURI host '" + host + "' cannot be 
resolved");
-            }
-            if (addr.isLoopbackAddress()
-                    || addr.isLinkLocalAddress()
-                    || addr.isSiteLocalAddress()
-                    || addr.isAnyLocalAddress()
-                    || addr.isMulticastAddress()) {
-                return Optional.of("jarURI host '" + host + "' resolves to a 
restricted address");
-            }
-        }
-        return Optional.empty();
+        return JarUriValidationUtils.validateJarURI(
+                jarURI, allowedSchemes, disallowRestrictedHosts);
     }
 
     private Optional<String> validateSessionJobJarURI(FlinkSessionJob 
sessionJob) {
diff --git 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/artifact/ArtifactManagerTest.java
 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/artifact/ArtifactManagerTest.java
index 422c08b1..1b785784 100644
--- 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/artifact/ArtifactManagerTest.java
+++ 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/artifact/ArtifactManagerTest.java
@@ -42,6 +42,7 @@ import java.net.HttpURLConnection;
 import java.net.InetSocketAddress;
 import java.net.URL;
 import java.nio.file.Path;
+import java.util.List;
 import java.util.Map;
 
 /** Test for {@link ArtifactManager}. */
@@ -53,11 +54,21 @@ public class ArtifactManagerTest {
 
     @BeforeEach
     public void setup() {
+        // The test server binds to loopback, so the operator policy must 
permit http + loopback.
+        artifactManager = artifactManagerWithPolicy(List.of("http"), false);
+    }
+
+    private ArtifactManager artifactManagerWithPolicy(
+            List<String> allowedSchemes, boolean disallowRestrictedHosts) {
         Configuration configuration = new Configuration();
         configuration.setString(
                 
KubernetesOperatorConfigOptions.OPERATOR_USER_ARTIFACTS_BASE_DIR,
                 tempDir.toAbsolutePath().toString());
-        artifactManager = new ArtifactManager(new 
FlinkConfigManager(configuration));
+        
configuration.set(KubernetesOperatorConfigOptions.JAR_URI_ALLOWED_SCHEMES, 
allowedSchemes);
+        configuration.set(
+                
KubernetesOperatorConfigOptions.JAR_URI_DISALLOW_RESTRICTED_HOSTS,
+                disallowRestrictedHosts);
+        return new ArtifactManager(new FlinkConfigManager(configuration));
     }
 
     @Test
@@ -117,6 +128,160 @@ public class ArtifactManagerTest {
         }
     }
 
+    @Test
+    public void testHttpFetchFollowsRedirectToAllowedTarget() throws Exception 
{
+        HttpServer httpServer = null;
+        try {
+            httpServer = startHttpServer();
+            var port = httpServer.getAddress().getPort();
+            var sourceFile = mockTheJarFile();
+            httpServer.createContext("/download/file.jar", new 
DownloadFileHttpHandler(sourceFile));
+            httpServer.createContext(
+                    "/myjob.jar",
+                    new RedirectHttpHandler(
+                            
String.format("http://127.0.0.1:%d/download/file.jar";, port)));
+
+            var file =
+                    artifactManager.fetch(
+                            String.format("http://127.0.0.1:%d/myjob.jar";, 
port),
+                            new Configuration(),
+                            tempDir.toString());
+            Assertions.assertTrue(file.exists());
+            // Content comes from the redirect target, but the name from the 
original jarURI.
+            Assertions.assertEquals("myjob.jar", file.getName());
+            Assertions.assertEquals(sourceFile.length(), file.length());
+        } finally {
+            if (httpServer != null) {
+                httpServer.stop(0);
+            }
+        }
+    }
+
+    @Test
+    public void testHttpFetchBlocksRedirectToNonHttpScheme() throws Exception {
+        // An http fetch must only follow http(s) redirects. Here a 
JDK-recognized non-http scheme
+        // (ftp) is rejected cleanly.
+        HttpServer httpServer = null;
+        try {
+            httpServer = startHttpServer();
+            var port = httpServer.getAddress().getPort();
+            httpServer.createContext(
+                    "/redirect",
+                    new 
RedirectHttpHandler(String.format("ftp://127.0.0.1:%d/job.jar";, port)));
+
+            var ex =
+                    Assertions.assertThrows(
+                            IOException.class,
+                            () ->
+                                    artifactManager.fetch(
+                                            
String.format("http://127.0.0.1:%d/redirect";, port),
+                                            new Configuration(),
+                                            tempDir.toString()));
+            Assertions.assertTrue(ex.getMessage().contains("non-http(s) 
target"), ex.getMessage());
+        } finally {
+            if (httpServer != null) {
+                httpServer.stop(0);
+            }
+        }
+    }
+
+    @Test
+    public void testHttpFetchBlocksRedirectToFilesystemScheme() throws 
Exception {
+        // An http server redirecting to an s3/hdfs target (a Flink filesystem 
scheme, not a
+        // java.net URL protocol) must fail closed cleanly rather than with a 
raw error.
+        HttpServer httpServer = null;
+        try {
+            httpServer = startHttpServer();
+            var port = httpServer.getAddress().getPort();
+            httpServer.createContext("/redirect", new 
RedirectHttpHandler("s3://bucket/job.jar"));
+
+            var ex =
+                    Assertions.assertThrows(
+                            IOException.class,
+                            () ->
+                                    artifactManager.fetch(
+                                            
String.format("http://127.0.0.1:%d/redirect";, port),
+                                            new Configuration(),
+                                            tempDir.toString()));
+            Assertions.assertTrue(
+                    ex.getMessage().contains("Refusing to follow redirect"), 
ex.getMessage());
+        } finally {
+            if (httpServer != null) {
+                httpServer.stop(0);
+            }
+        }
+    }
+
+    @Test
+    public void testHttpFetchBlocksTooManyRedirects() throws Exception {
+        HttpServer httpServer = null;
+        try {
+            httpServer = startHttpServer();
+            var port = httpServer.getAddress().getPort();
+            httpServer.createContext(
+                    "/loop",
+                    new 
RedirectHttpHandler(String.format("http://127.0.0.1:%d/loop";, port)));
+
+            var ex =
+                    Assertions.assertThrows(
+                            IOException.class,
+                            () ->
+                                    artifactManager.fetch(
+                                            
String.format("http://127.0.0.1:%d/loop";, port),
+                                            new Configuration(),
+                                            tempDir.toString()));
+            Assertions.assertTrue(ex.getMessage().contains("Too many 
redirects"), ex.getMessage());
+        } finally {
+            if (httpServer != null) {
+                httpServer.stop(0);
+            }
+        }
+    }
+
+    @Test
+    public void testOperatorConfigControlsRestrictedHostPolicy() {
+        // The restricted-host policy comes from the operator config; a value 
set in the per-job
+        // config does not override it. No server is needed: the loopback host 
is rejected first.
+        var strictManager = artifactManagerWithPolicy(List.of("http"), true);
+        var jobConfig =
+                new Configuration()
+                        .set(
+                                
KubernetesOperatorConfigOptions.JAR_URI_DISALLOW_RESTRICTED_HOSTS,
+                                false);
+
+        var ex =
+                Assertions.assertThrows(
+                        IOException.class,
+                        () ->
+                                strictManager.fetch(
+                                        "http://127.0.0.1:9999/job.jar";,
+                                        jobConfig,
+                                        tempDir.toString()));
+        Assertions.assertTrue(ex.getMessage().contains("restricted address"), 
ex.getMessage());
+    }
+
+    @Test
+    public void testOperatorConfigControlsSchemeAllowlist() {
+        // The scheme allowlist comes from the operator config; a value set in 
the per-job config
+        // does not override it.
+        var strictManager = artifactManagerWithPolicy(List.of("https"), false);
+        var jobConfig =
+                new Configuration()
+                        .set(
+                                
KubernetesOperatorConfigOptions.JAR_URI_ALLOWED_SCHEMES,
+                                List.of("http"));
+
+        var ex =
+                Assertions.assertThrows(
+                        IOException.class,
+                        () ->
+                                strictManager.fetch(
+                                        "http://127.0.0.1:9999/job.jar";,
+                                        jobConfig,
+                                        tempDir.toString()));
+        Assertions.assertTrue(ex.getMessage().contains("scheme 'http'"), 
ex.getMessage());
+    }
+
     private HttpServer startHttpServer() throws IOException {
         int port = RandomUtils.nextInt(2000, 3000);
         HttpServer httpServer = null;
@@ -160,4 +325,21 @@ public class ArtifactManagerTest {
             exchange.close();
         }
     }
+
+    /** Handler that always responds with a 302 redirect to the configured 
location. */
+    public static class RedirectHttpHandler implements HttpHandler {
+
+        private final String location;
+
+        public RedirectHttpHandler(String location) {
+            this.location = location;
+        }
+
+        @Override
+        public void handle(HttpExchange exchange) throws IOException {
+            exchange.getResponseHeaders().add("Location", location);
+            exchange.sendResponseHeaders(HttpURLConnection.HTTP_MOVED_TEMP, 
-1);
+            exchange.close();
+        }
+    }
 }

Reply via email to