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 91f543f0 [FLINK-40402] Harden job artifact fetching against slow-host
reconcile-thread exhaustion (#1186)
91f543f0 is described below
commit 91f543f08feec2fcd0d64c0c5e086cd075836046
Author: Purushottam Sinha <[email protected]>
AuthorDate: Fri Aug 21 13:17:18 2026 +0530
[FLINK-40402] Harden job artifact fetching against slow-host
reconcile-thread exhaustion (#1186)
* [FLINK-40402] Harden job artifact fetching against slow-host
reconcile-thread exhaustion.
Add connect/read/total timeouts and a size cap to HttpArtifactFetcher so a
slow or unresponsive FlinkSessionJob jarURI host can no longer pin a
reconcile thread indefinitely.
Generated-by: Claude Code
* [FLINK-40402] Make artifact-fetch failure messages consistent and
traceable.
The size-cap and total-timeout errors thrown while streaming the body did
not
name the offending jarURI or the configured value, unlike the declared-size
and redirect-timeout errors thrown earlier in the fetch. Pass the URI and
the
total timeout into copyBounded so every fetch failure names the artifact,
and
unify the wording: both size errors now read "exceeds the configured limit"
and both timeout errors read "Timed out (> <duration>) while fetching
artifact
from '<uri>'".
Generated-by: Claude Code
* [FLINK-40402] Add regression tests for the connect and read fetch
timeouts.
The existing coverage only exercised the overall (total) fetch timeout via a
slow trickle. Add two tests that isolate the per-connection socket timeouts,
each with a much longer total timeout so it is unambiguously the socket
timeout that bounds them:
- connect timeout: fetching from a blackholed address (TEST-NET-1) whose
SYNs are dropped, so the TCP connect never completes;
- read timeout: a server that sends response headers then stalls without
sending the body, so a client read blocks.
Generated-by: Claude Code
---
docs/content.zh/docs/deployment/security.md | 6 +
docs/content/docs/deployment/security.md | 6 +
.../kubernetes_operator_config_configuration.html | 18 ++
.../generated/system_reconcile_section.html | 18 ++
.../operator/artifact/ArtifactManager.java | 16 +-
.../operator/artifact/HttpArtifactFetcher.java | 114 +++++++-
.../config/FlinkOperatorConfiguration.java | 18 +-
.../config/KubernetesOperatorConfigOptions.java | 34 +++
.../operator/artifact/ArtifactManagerTest.java | 311 +++++++++++++++++++++
9 files changed, 515 insertions(+), 26 deletions(-)
diff --git a/docs/content.zh/docs/deployment/security.md
b/docs/content.zh/docs/deployment/security.md
index e1dd778d..9236604b 100644
--- a/docs/content.zh/docs/deployment/security.md
+++ b/docs/content.zh/docs/deployment/security.md
@@ -66,6 +66,12 @@ For a `FlinkSessionJob`, the operator itself downloads the
job artifact referenc
- `kubernetes.operator.user.artifacts.disallow-restricted-hosts` (default
`true`): rejects `http` and `https` URIs whose host resolves to a loopback,
link-local, site-local, wildcard, or multicast address, so a session job cannot
point the operator at cluster-internal endpoints.
- `kubernetes.operator.user.artifacts.http.header`: custom HTTP headers sent
when fetching artifacts over `http` and `https`, typically carrying the
credentials of the artifact store.
+The fetch itself runs synchronously on the reconcile thread, so a slow or
unresponsive artifact host could otherwise pin that thread and, with enough
concurrent session jobs, exhaust the bounded reconcile pool
(`kubernetes.operator.reconcile.parallelism`). Two settings bound the reconcile
thread's exposure to such a host, and one bounds the response size:
+
+- `kubernetes.operator.user.artifacts.http.socket-timeout` (default `30 s`):
the connect and per-read socket timeout for the underlying HTTP connection.
+- `kubernetes.operator.user.artifacts.http.total-timeout` (default `5 min`):
the overall wall-clock budget for the whole fetch, covering all redirects and
the full body transfer. This is what bounds a host that trickles data just
slowly enough to keep individual reads under the socket timeout without the
transfer ever completing.
+- `kubernetes.operator.user.artifacts.max-size` (default `1 gb`): the maximum
artifact size. The download is rejected once it exceeds this, whether or not
the server declares a `Content-Length` up front.
+
Application-mode deployments are not affected: their `jarURI` is resolved
inside the job's own cluster, not by the operator.
## Secrets in Configuration
diff --git a/docs/content/docs/deployment/security.md
b/docs/content/docs/deployment/security.md
index e1dd778d..9236604b 100644
--- a/docs/content/docs/deployment/security.md
+++ b/docs/content/docs/deployment/security.md
@@ -66,6 +66,12 @@ For a `FlinkSessionJob`, the operator itself downloads the
job artifact referenc
- `kubernetes.operator.user.artifacts.disallow-restricted-hosts` (default
`true`): rejects `http` and `https` URIs whose host resolves to a loopback,
link-local, site-local, wildcard, or multicast address, so a session job cannot
point the operator at cluster-internal endpoints.
- `kubernetes.operator.user.artifacts.http.header`: custom HTTP headers sent
when fetching artifacts over `http` and `https`, typically carrying the
credentials of the artifact store.
+The fetch itself runs synchronously on the reconcile thread, so a slow or
unresponsive artifact host could otherwise pin that thread and, with enough
concurrent session jobs, exhaust the bounded reconcile pool
(`kubernetes.operator.reconcile.parallelism`). Two settings bound the reconcile
thread's exposure to such a host, and one bounds the response size:
+
+- `kubernetes.operator.user.artifacts.http.socket-timeout` (default `30 s`):
the connect and per-read socket timeout for the underlying HTTP connection.
+- `kubernetes.operator.user.artifacts.http.total-timeout` (default `5 min`):
the overall wall-clock budget for the whole fetch, covering all redirects and
the full body transfer. This is what bounds a host that trickles data just
slowly enough to keep individual reads under the socket timeout without the
transfer ever completing.
+- `kubernetes.operator.user.artifacts.max-size` (default `1 gb`): the maximum
artifact size. The download is rejected once it exceeds this, whether or not
the server declares a `Content-Length` up front.
+
Application-mode deployments are not affected: their `jarURI` is resolved
inside the job's own cluster, not by the operator.
## Secrets in Configuration
diff --git
a/docs/layouts/shortcodes/generated/kubernetes_operator_config_configuration.html
b/docs/layouts/shortcodes/generated/kubernetes_operator_config_configuration.html
index cd8513bf..502cedb2 100644
---
a/docs/layouts/shortcodes/generated/kubernetes_operator_config_configuration.html
+++
b/docs/layouts/shortcodes/generated/kubernetes_operator_config_configuration.html
@@ -500,6 +500,24 @@
<td>Map</td>
<td>Custom HTTP header for HttpArtifactFetcher. The header will be
applied when getting the session job artifacts. Expected format:
headerKey1:headerValue1,headerKey2:headerValue2.</td>
</tr>
+ <tr>
+
<td><h5>kubernetes.operator.user.artifacts.http.socket-timeout</h5></td>
+ <td style="word-wrap: break-word;">30 s</td>
+ <td>Duration</td>
+ <td>The connect and socket read timeout for downloading a
FlinkSessionJob jarURI over http(s). Bounds how long the reconcile thread can
be blocked establishing the connection, or waiting for the next byte, from an
unresponsive artifact host.</td>
+ </tr>
+ <tr>
+
<td><h5>kubernetes.operator.user.artifacts.http.total-timeout</h5></td>
+ <td style="word-wrap: break-word;">5 min</td>
+ <td>Duration</td>
+ <td>The total wall-clock budget for downloading a FlinkSessionJob
jarURI over http(s), covering all redirects and the full body transfer. Unlike
the socket timeout, this bounds the overall download even against a slow host
that keeps trickling data fast enough to avoid tripping it.</td>
+ </tr>
+ <tr>
+ <td><h5>kubernetes.operator.user.artifacts.max-size</h5></td>
+ <td style="word-wrap: break-word;">1 gb</td>
+ <td>MemorySize</td>
+ <td>The maximum size of a FlinkSessionJob jarURI artifact fetched
over http(s). The download is rejected once it exceeds this size, whether or
not the server declares a Content-Length up front.</td>
+ </tr>
<tr>
<td><h5>kubernetes.operator.watched.namespaces</h5></td>
<td style="word-wrap: break-word;">"JOSDK_ALL_NAMESPACES"</td>
diff --git a/docs/layouts/shortcodes/generated/system_reconcile_section.html
b/docs/layouts/shortcodes/generated/system_reconcile_section.html
index bb2769c4..5633ace9 100644
--- a/docs/layouts/shortcodes/generated/system_reconcile_section.html
+++ b/docs/layouts/shortcodes/generated/system_reconcile_section.html
@@ -80,5 +80,23 @@
<td>Boolean</td>
<td>If enabled, FlinkSessionJob jarURI hosts that resolve to
loopback, link-local, site-local, wildcard or multicast addresses are rejected
during validation. Disable only if the operator legitimately needs to fetch
from such addresses.</td>
</tr>
+ <tr>
+
<td><h5>kubernetes.operator.user.artifacts.http.socket-timeout</h5></td>
+ <td style="word-wrap: break-word;">30 s</td>
+ <td>Duration</td>
+ <td>The connect and socket read timeout for downloading a
FlinkSessionJob jarURI over http(s). Bounds how long the reconcile thread can
be blocked establishing the connection, or waiting for the next byte, from an
unresponsive artifact host.</td>
+ </tr>
+ <tr>
+
<td><h5>kubernetes.operator.user.artifacts.http.total-timeout</h5></td>
+ <td style="word-wrap: break-word;">5 min</td>
+ <td>Duration</td>
+ <td>The total wall-clock budget for downloading a FlinkSessionJob
jarURI over http(s), covering all redirects and the full body transfer. Unlike
the socket timeout, this bounds the overall download even against a slow host
that keeps trickling data fast enough to avoid tripping it.</td>
+ </tr>
+ <tr>
+ <td><h5>kubernetes.operator.user.artifacts.max-size</h5></td>
+ <td style="word-wrap: break-word;">1 gb</td>
+ <td>MemorySize</td>
+ <td>The maximum size of a FlinkSessionJob jarURI artifact fetched
over http(s). The download is rejected once it exceeds this size, whether or
not the server declares a Content-Length up front.</td>
+ </tr>
</tbody>
</table>
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 7e229573..d755a14b 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,7 +20,6 @@ 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;
@@ -59,17 +58,12 @@ public class ArtifactManager {
createIfNotExists(targetDir);
URI uri = new URI(jarURI);
if ("http".equals(uri.getScheme()) || "https".equals(uri.getScheme()))
{
- // Take the scheme/host policy from the operator config (matching
DefaultValidator);
- // clone so the caller's config is not mutated.
+ // The scheme/host policy, fetch timeouts and size cap come from
the operator config
+ // (matching DefaultValidator), not the tenant-influenced
flinkConfiguration, so it's
+ // passed to the fetcher separately rather than merged into
flinkConfiguration.
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);
+ return HttpArtifactFetcher.INSTANCE.fetch(
+ jarURI, flinkConfiguration, operatorConfig, 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 09e517c6..3a7a8b74 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
@@ -18,6 +18,7 @@
package org.apache.flink.kubernetes.operator.artifact;
import org.apache.flink.configuration.Configuration;
+import org.apache.flink.kubernetes.operator.config.FlinkOperatorConfiguration;
import
org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions;
import org.apache.flink.kubernetes.operator.utils.JarUriValidationUtils;
@@ -27,18 +28,22 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.File;
+import java.io.FileOutputStream;
import java.io.IOException;
+import java.io.InputStream;
import java.net.HttpURLConnection;
import java.net.MalformedURLException;
import java.net.URL;
+import java.time.Duration;
import java.util.Map;
/**
- * 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.
+ * Download the jar from the http resource. The scheme allowlist,
restricted-host policy, fetch
+ * timeouts and size cap come from the trusted operator configuration passed
to {@link #fetch}, not
+ * the (possibly tenant-influenced) {@code flinkConfiguration}, which only
supplies the HTTP
+ * headers.
*/
-public class HttpArtifactFetcher implements ArtifactFetcher {
+public class HttpArtifactFetcher {
public static final Logger LOG =
LoggerFactory.getLogger(HttpArtifactFetcher.class);
public static final HttpArtifactFetcher INSTANCE = new
HttpArtifactFetcher();
@@ -46,18 +51,26 @@ public class HttpArtifactFetcher implements ArtifactFetcher
{
// 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)
+ // Chunk size used when streaming the response body to disk.
+ private static final int COPY_BUFFER_SIZE = 8 * 1024;
+
+ public File fetch(
+ String uri,
+ Configuration flinkConfiguration,
+ FlinkOperatorConfiguration operatorConfig,
+ File targetDir)
throws Exception {
var start = System.currentTimeMillis();
- // 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);
+ var allowedSchemes = operatorConfig.getJarUriAllowedSchemes();
+ var disallowRestrictedHosts =
operatorConfig.isJarUriDisallowRestrictedHosts();
+ var socketTimeoutMillis = (int)
operatorConfig.getJarFetchSocketTimeout().toMillis();
+ var totalTimeout = operatorConfig.getJarFetchTotalTimeout();
+ var maxArtifactSize =
operatorConfig.getJarArtifactMaxSize().getBytes();
+ // Overall wall-clock deadline for the whole fetch (all redirects +
the body transfer).
+ // This bounds the reconcile thread even against a host that keeps
trickling data slowly
+ // enough to never trip the socket timeout on its own.
+ var deadline = start + totalTimeout.toMillis();
// merged session job level header and cluster level header, session
job level header take
// precedence.
@@ -72,6 +85,14 @@ public class HttpArtifactFetcher implements ArtifactFetcher {
HttpURLConnection conn;
int redirects = 0;
while (true) {
+ if (System.currentTimeMillis() > deadline) {
+ throw new IOException(
+ "Timed out (> "
+ + totalTimeout
+ + ") while fetching artifact from '"
+ + uri
+ + "'");
+ }
var validationError =
JarUriValidationUtils.validateJarURI(
currentUri, allowedSchemes,
disallowRestrictedHosts);
@@ -89,6 +110,8 @@ public class HttpArtifactFetcher implements ArtifactFetcher {
}
conn = (HttpURLConnection) currentUrl.openConnection();
conn.setInstanceFollowRedirects(false);
+ conn.setConnectTimeout(socketTimeoutMillis);
+ conn.setReadTimeout(socketTimeoutMillis);
// Only send the configured headers to the original host; drop
them on a cross-host
// redirect.
if (headers != null &&
originalUrl.getHost().equalsIgnoreCase(currentUrl.getHost())) {
@@ -154,12 +177,30 @@ public class HttpArtifactFetcher implements
ArtifactFetcher {
}
}
+ // Fail fast if the server declares a size beyond the cap; a
malicious/misconfigured
+ // server can still lie about this, so the copy below enforces the cap
regardless.
+ long declaredLength = conn.getContentLengthLong();
+ if (declaredLength > maxArtifactSize) {
+ conn.disconnect();
+ throw new IOException(
+ "Refusing to fetch artifact from '"
+ + uri
+ + "': declared size "
+ + declaredLength
+ + " bytes exceeds the configured limit of "
+ + maxArtifactSize
+ + " bytes");
+ }
+
// 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);
+ copyBounded(inputStream, targetFile, maxArtifactSize, deadline,
uri, totalTimeout);
+ } catch (Exception e) {
+ targetFile.delete();
+ throw e;
} finally {
conn.disconnect();
}
@@ -171,6 +212,51 @@ public class HttpArtifactFetcher implements
ArtifactFetcher {
return targetFile;
}
+ /**
+ * Streams {@code inputStream} to {@code targetFile}, aborting if the
total bytes written exceed
+ * {@code maxBytes} or {@code deadline} (wall-clock millis) passes. Each
individual {@link
+ * InputStream#read} is already bounded by the connection's read timeout,
so the deadline check
+ * here is what catches a host that trickles data just fast enough to keep
each read below that
+ * timeout without ever finishing.
+ */
+ private static void copyBounded(
+ InputStream inputStream,
+ File targetFile,
+ long maxBytes,
+ long deadline,
+ String uri,
+ Duration totalTimeout)
+ throws IOException {
+ FileUtils.forceMkdirParent(targetFile);
+ byte[] buffer = new byte[COPY_BUFFER_SIZE];
+ long total = 0;
+ try (var outputStream = new FileOutputStream(targetFile)) {
+ int read;
+ while ((read = inputStream.read(buffer)) != -1) {
+ total += read;
+ if (total > maxBytes) {
+ throw new IOException(
+ "Refusing to fetch artifact from '"
+ + uri
+ + "': downloaded size "
+ + total
+ + " bytes exceeds the configured limit of "
+ + maxBytes
+ + " bytes");
+ }
+ if (System.currentTimeMillis() > deadline) {
+ throw new IOException(
+ "Timed out (> "
+ + totalTimeout
+ + ") while fetching artifact from '"
+ + uri
+ + "'");
+ }
+ outputStream.write(buffer, 0, read);
+ }
+ }
+ }
+
private static boolean isRedirect(int status) {
return status == HttpURLConnection.HTTP_MOVED_PERM
|| status == HttpURLConnection.HTTP_MOVED_TEMP
diff --git
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/FlinkOperatorConfiguration.java
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/FlinkOperatorConfiguration.java
index 14da8747..f70198f6 100644
---
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/FlinkOperatorConfiguration.java
+++
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/FlinkOperatorConfiguration.java
@@ -20,6 +20,7 @@ package org.apache.flink.kubernetes.operator.config;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.configuration.IllegalConfigurationException;
+import org.apache.flink.configuration.MemorySize;
import
org.apache.flink.kubernetes.operator.metrics.KubernetesOperatorMetricOptions;
import org.apache.flink.kubernetes.operator.utils.EnvUtils;
@@ -84,6 +85,9 @@ public class FlinkOperatorConfiguration {
Duration jobSubmissionTimeout;
List<String> jarUriAllowedSchemes;
boolean jarUriDisallowRestrictedHosts;
+ Duration jarFetchSocketTimeout;
+ Duration jarFetchTotalTimeout;
+ MemorySize jarArtifactMaxSize;
public static FlinkOperatorConfiguration fromConfiguration(Configuration
operatorConfig) {
Duration reconcileInterval =
@@ -221,6 +225,15 @@ public class FlinkOperatorConfiguration {
operatorConfig.get(
KubernetesOperatorConfigOptions.JAR_URI_DISALLOW_RESTRICTED_HOSTS);
+ Duration jarFetchSocketTimeout =
+
operatorConfig.get(KubernetesOperatorConfigOptions.JAR_FETCH_SOCKET_TIMEOUT);
+
+ Duration jarFetchTotalTimeout =
+
operatorConfig.get(KubernetesOperatorConfigOptions.JAR_FETCH_TOTAL_TIMEOUT);
+
+ MemorySize jarArtifactMaxSize =
+
operatorConfig.get(KubernetesOperatorConfigOptions.JAR_ARTIFACT_MAX_SIZE);
+
return new FlinkOperatorConfiguration(
reconcileInterval,
reconcilerMaxParallelism,
@@ -256,7 +269,10 @@ public class FlinkOperatorConfiguration {
manageIngress,
jobSubmissionTimeout,
jarUriAllowedSchemes,
- jarUriDisallowRestrictedHosts);
+ jarUriDisallowRestrictedHosts,
+ jarFetchSocketTimeout,
+ jarFetchTotalTimeout,
+ jarArtifactMaxSize);
}
private static GenericRetry getRetryConfig(Configuration conf) {
diff --git
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/KubernetesOperatorConfigOptions.java
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/KubernetesOperatorConfigOptions.java
index 3a4d10b4..fb98eba9 100644
---
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/KubernetesOperatorConfigOptions.java
+++
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/KubernetesOperatorConfigOptions.java
@@ -21,6 +21,7 @@ package org.apache.flink.kubernetes.operator.config;
import org.apache.flink.annotation.docs.Documentation;
import org.apache.flink.configuration.ConfigOption;
import org.apache.flink.configuration.ConfigOptions;
+import org.apache.flink.configuration.MemorySize;
import org.apache.flink.core.execution.SavepointFormatType;
import org.apache.flink.kubernetes.operator.api.status.CheckpointType;
@@ -397,6 +398,39 @@ public class KubernetesOperatorConfigOptions {
+ "site-local, wildcard or multicast
addresses are rejected during validation. "
+ "Disable only if the operator
legitimately needs to fetch from such addresses.");
+ @Documentation.Section(SECTION_SYSTEM_RECONCILE)
+ public static final ConfigOption<Duration> JAR_FETCH_SOCKET_TIMEOUT =
+ operatorConfig("user.artifacts.http.socket-timeout")
+ .durationType()
+ .defaultValue(Duration.ofSeconds(30))
+ .withDescription(
+ "The connect and socket read timeout for
downloading a FlinkSessionJob jarURI "
+ + "over http(s). Bounds how long the
reconcile thread can be blocked "
+ + "establishing the connection, or waiting
for the next byte, from an "
+ + "unresponsive artifact host.");
+
+ @Documentation.Section(SECTION_SYSTEM_RECONCILE)
+ public static final ConfigOption<Duration> JAR_FETCH_TOTAL_TIMEOUT =
+ operatorConfig("user.artifacts.http.total-timeout")
+ .durationType()
+ .defaultValue(Duration.ofMinutes(5))
+ .withDescription(
+ "The total wall-clock budget for downloading a
FlinkSessionJob jarURI over "
+ + "http(s), covering all redirects and the
full body transfer. Unlike "
+ + "the socket timeout, this bounds the
overall download even against "
+ + "a slow host that keeps trickling data
fast enough to avoid tripping "
+ + "it.");
+
+ @Documentation.Section(SECTION_SYSTEM_RECONCILE)
+ public static final ConfigOption<MemorySize> JAR_ARTIFACT_MAX_SIZE =
+ operatorConfig("user.artifacts.max-size")
+ .memoryType()
+ .defaultValue(MemorySize.ofMebiBytes(1024))
+ .withDescription(
+ "The maximum size of a FlinkSessionJob jarURI
artifact fetched over http(s). "
+ + "The download is rejected once it
exceeds this size, whether or not the "
+ + "server declares a Content-Length up
front.");
+
@Documentation.Section(SECTION_DYNAMIC)
public static final ConfigOption<Boolean> SNAPSHOT_RESOURCE_ENABLED =
operatorConfig("snapshot.resource.enabled")
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 1b785784..87a6f8e7 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
@@ -18,6 +18,7 @@
package org.apache.flink.kubernetes.operator.artifact;
import org.apache.flink.configuration.Configuration;
+import org.apache.flink.configuration.MemorySize;
import org.apache.flink.kubernetes.operator.TestUtils;
import org.apache.flink.kubernetes.operator.config.FlinkConfigManager;
import
org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions;
@@ -42,6 +43,7 @@ import java.net.HttpURLConnection;
import java.net.InetSocketAddress;
import java.net.URL;
import java.nio.file.Path;
+import java.time.Duration;
import java.util.List;
import java.util.Map;
@@ -71,6 +73,34 @@ public class ArtifactManagerTest {
return new ArtifactManager(new FlinkConfigManager(configuration));
}
+ private ArtifactManager artifactManagerWithFetchLimits(
+ Duration totalTimeout, long maxArtifactSizeBytes) {
+ Configuration configuration = new Configuration();
+ configuration.setString(
+
KubernetesOperatorConfigOptions.OPERATOR_USER_ARTIFACTS_BASE_DIR,
+ tempDir.toAbsolutePath().toString());
+
configuration.set(KubernetesOperatorConfigOptions.JAR_URI_ALLOWED_SCHEMES,
List.of("http"));
+
configuration.set(KubernetesOperatorConfigOptions.JAR_URI_DISALLOW_RESTRICTED_HOSTS,
false);
+
configuration.set(KubernetesOperatorConfigOptions.JAR_FETCH_TOTAL_TIMEOUT,
totalTimeout);
+ configuration.set(
+ KubernetesOperatorConfigOptions.JAR_ARTIFACT_MAX_SIZE,
+ new MemorySize(maxArtifactSizeBytes));
+ return new ArtifactManager(new FlinkConfigManager(configuration));
+ }
+
+ private ArtifactManager artifactManagerWithSocketTimeout(
+ Duration socketTimeout, Duration totalTimeout) {
+ Configuration configuration = new Configuration();
+ configuration.setString(
+
KubernetesOperatorConfigOptions.OPERATOR_USER_ARTIFACTS_BASE_DIR,
+ tempDir.toAbsolutePath().toString());
+
configuration.set(KubernetesOperatorConfigOptions.JAR_URI_ALLOWED_SCHEMES,
List.of("http"));
+
configuration.set(KubernetesOperatorConfigOptions.JAR_URI_DISALLOW_RESTRICTED_HOSTS,
false);
+
configuration.set(KubernetesOperatorConfigOptions.JAR_FETCH_SOCKET_TIMEOUT,
socketTimeout);
+
configuration.set(KubernetesOperatorConfigOptions.JAR_FETCH_TOTAL_TIMEOUT,
totalTimeout);
+ return new ArtifactManager(new FlinkConfigManager(configuration));
+ }
+
@Test
public void testGenerateJarDir() {
var sessionJob = TestUtils.buildSessionJob();
@@ -128,6 +158,36 @@ public class ArtifactManagerTest {
}
}
+ @Test
+ public void testHttpFetchCreatesNestedNonExistentTargetDir() throws
Exception {
+ // The real call site (uploadJar -> generateJarDir) targets a nested
per-job directory
+ // (base/namespace/deployment/job) that doesn't exist yet; unlike the
other tests here,
+ // don't reuse the JUnit-provided tempDir directly so a missing
intermediate directory
+ // actually gets exercised.
+ var nestedTargetDir =
tempDir.resolve("ns").resolve("deployment").resolve("job");
+ Assertions.assertFalse(nestedTargetDir.toFile().exists());
+ HttpServer httpServer = null;
+ try {
+ httpServer = startHttpServer();
+ var sourceFile = mockTheJarFile();
+ httpServer.createContext("/download/file.jar", new
DownloadFileHttpHandler(sourceFile));
+
+ var file =
+ artifactManager.fetch(
+ String.format(
+ "http://127.0.0.1:%d/download/file.jar",
+ httpServer.getAddress().getPort()),
+ new Configuration(),
+ nestedTargetDir.toString());
+ Assertions.assertTrue(file.exists());
+ Assertions.assertEquals(nestedTargetDir.toString(),
file.getParent());
+ } finally {
+ if (httpServer != null) {
+ httpServer.stop(0);
+ }
+ }
+ }
+
@Test
public void testHttpFetchFollowsRedirectToAllowedTarget() throws Exception
{
HttpServer httpServer = null;
@@ -238,6 +298,191 @@ public class ArtifactManagerTest {
}
}
+ @Test
+ public void testHttpFetchRejectsOversizedDeclaredContentLength() throws
Exception {
+ // The server declares a Content-Length beyond the cap; the fetch must
fail fast without
+ // reading the body.
+ var strictManager =
artifactManagerWithFetchLimits(Duration.ofSeconds(30), 10);
+ HttpServer httpServer = null;
+ try {
+ httpServer = startHttpServer();
+ var port = httpServer.getAddress().getPort();
+ var sourceFile = mockTheJarFile();
+ Assertions.assertTrue(sourceFile.length() > 10);
+ httpServer.createContext("/download/file.jar", new
DownloadFileHttpHandler(sourceFile));
+
+ var ex =
+ Assertions.assertThrows(
+ IOException.class,
+ () ->
+ strictManager.fetch(
+ String.format(
+
"http://127.0.0.1:%d/download/file.jar", port),
+ new Configuration(),
+ tempDir.toString()));
+ Assertions.assertTrue(ex.getMessage().contains("exceeds the
configured limit"));
+ } finally {
+ if (httpServer != null) {
+ httpServer.stop(0);
+ }
+ }
+ }
+
+ @Test
+ public void testHttpFetchRejectsOversizedActualBody() throws Exception {
+ // The server does not declare a Content-Length (chunked transfer), so
the cap must be
+ // enforced while streaming the body rather than up front.
+ var strictManager =
artifactManagerWithFetchLimits(Duration.ofSeconds(30), 10);
+ HttpServer httpServer = null;
+ try {
+ httpServer = startHttpServer();
+ var port = httpServer.getAddress().getPort();
+ httpServer.createContext("/download/file.jar", new
ChunkedOversizedHttpHandler());
+
+ var ex =
+ Assertions.assertThrows(
+ IOException.class,
+ () ->
+ strictManager.fetch(
+ String.format(
+
"http://127.0.0.1:%d/download/file.jar", port),
+ new Configuration(),
+ tempDir.toString()));
+ Assertions.assertTrue(ex.getMessage().contains("exceeds the
configured limit"));
+ Assertions.assertTrue(ex.getMessage().contains("downloaded size"),
ex.getMessage());
+ Assertions.assertFalse(new File(tempDir.toFile(),
"file.jar").exists());
+ } finally {
+ if (httpServer != null) {
+ httpServer.stop(0);
+ }
+ }
+ }
+
+ @Test
+ public void testHttpFetchTimesOutOnSlowTrickle() throws Exception {
+ // The server sends a byte at a time, each well within the read
timeout, so only the
+ // overall fetch timeout can bound the reconcile thread here.
+ var strictManager =
artifactManagerWithFetchLimits(Duration.ofMillis(300), 10_000_000);
+ HttpServer httpServer = null;
+ try {
+ httpServer = startHttpServer();
+ var port = httpServer.getAddress().getPort();
+ httpServer.createContext("/download/file.jar", new
SlowTrickleHttpHandler());
+
+ var fetchStart = System.currentTimeMillis();
+ var ex =
+ Assertions.assertThrows(
+ IOException.class,
+ () ->
+ strictManager.fetch(
+ String.format(
+
"http://127.0.0.1:%d/download/file.jar", port),
+ new Configuration(),
+ tempDir.toString()));
+ var elapsed = System.currentTimeMillis() - fetchStart;
+ Assertions.assertTrue(ex.getMessage().contains("Timed out"),
ex.getMessage());
+ // Bounded well below the many seconds the slow trickle would
otherwise take to finish.
+ Assertions.assertTrue(elapsed < 10_000, "fetch took " + elapsed +
" ms");
+ } finally {
+ if (httpServer != null) {
+ httpServer.stop(0);
+ }
+ }
+ }
+
+ @Test
+ public void testHttpFetchTimesOutOnConnectToUnreachableHost() {
+ // Connect to a blackholed address (TEST-NET-1, RFC 5737) whose SYNs
are dropped, so the
+ // TCP connect never completes. The 1s connect timeout must bound it
far below the 60s
+ // total timeout, proving the per-connection timeout (not just the
total) is in effect.
+ var strictManager =
+ artifactManagerWithSocketTimeout(Duration.ofSeconds(1),
Duration.ofSeconds(60));
+
+ var fetchStart = System.currentTimeMillis();
+ var ex =
+ Assertions.assertThrows(
+ IOException.class,
+ () ->
+ strictManager.fetch(
+ "http://192.0.2.1:81/job.jar",
+ new Configuration(),
+ tempDir.toString()));
+ var elapsed = System.currentTimeMillis() - fetchStart;
+ Assertions.assertTrue(
+ ex.getMessage().toLowerCase().contains("connect timed out"),
ex.getMessage());
+ Assertions.assertTrue(elapsed < 30_000, "fetch took " + elapsed + "
ms");
+ }
+
+ @Test
+ public void testHttpFetchTimesOutOnReadWhenServerStallsAfterHeaders()
throws Exception {
+ // The server sends response headers then never sends the body, so a
client read blocks.
+ // The 1s read timeout must bound the stalled read far below the 60s
total timeout,
+ // proving the per-read socket timeout (not just the total) is in
effect.
+ var strictManager =
+ artifactManagerWithSocketTimeout(Duration.ofSeconds(1),
Duration.ofSeconds(60));
+ HttpServer httpServer = null;
+ try {
+ httpServer = startHttpServer();
+ var port = httpServer.getAddress().getPort();
+ httpServer.createContext("/download/file.jar", new
StallAfterHeadersHttpHandler());
+
+ var fetchStart = System.currentTimeMillis();
+ var ex =
+ Assertions.assertThrows(
+ IOException.class,
+ () ->
+ strictManager.fetch(
+ String.format(
+
"http://127.0.0.1:%d/download/file.jar", port),
+ new Configuration(),
+ tempDir.toString()));
+ var elapsed = System.currentTimeMillis() - fetchStart;
+ Assertions.assertTrue(
+ ex.getMessage().toLowerCase().contains("read timed out"),
ex.getMessage());
+ Assertions.assertTrue(elapsed < 30_000, "fetch took " + elapsed +
" ms");
+ } finally {
+ if (httpServer != null) {
+ httpServer.stop(0);
+ }
+ }
+ }
+
+ @Test
+ public void testOperatorConfigControlsFetchLimits() throws Exception {
+ // The size-cap policy comes from the operator config; a value set in
the per-job config
+ // does not override it.
+ var strictManager =
artifactManagerWithFetchLimits(Duration.ofSeconds(30), 10);
+ var jobConfig =
+ new Configuration()
+ .set(
+
KubernetesOperatorConfigOptions.JAR_ARTIFACT_MAX_SIZE,
+ MemorySize.ofMebiBytes(1024));
+
+ HttpServer httpServer = null;
+ try {
+ httpServer = startHttpServer();
+ var port = httpServer.getAddress().getPort();
+ var sourceFile = mockTheJarFile();
+ Assertions.assertTrue(sourceFile.length() > 10);
+ httpServer.createContext("/download/file.jar", new
DownloadFileHttpHandler(sourceFile));
+
+ var ex =
+ Assertions.assertThrows(
+ IOException.class,
+ () ->
+ strictManager.fetch(
+ String.format(
+
"http://127.0.0.1:%d/download/file.jar", port),
+ jobConfig,
+ tempDir.toString()));
+ Assertions.assertTrue(ex.getMessage().contains("exceeds the
configured limit"));
+ } 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
@@ -342,4 +587,70 @@ public class ArtifactManagerTest {
exchange.close();
}
}
+
+ /**
+ * Handler that streams more bytes than any reasonable test cap without
ever declaring a
+ * Content-Length (chunked transfer), so the size cap must be enforced
while streaming.
+ */
+ public static class ChunkedOversizedHttpHandler implements HttpHandler {
+
+ @Override
+ public void handle(HttpExchange exchange) throws IOException {
+ // -1 body length tells the JDK HTTP server to use chunked
transfer encoding.
+ exchange.sendResponseHeaders(HttpURLConnection.HTTP_OK, 0);
+ var body = exchange.getResponseBody();
+ byte[] chunk = new byte[1024];
+ for (int i = 0; i < 100; i++) {
+ body.write(chunk);
+ }
+ exchange.close();
+ }
+ }
+
+ /**
+ * Handler that dribbles a single byte at a time with a short pause
between each, simulating a
+ * slow-loris style host: each individual read completes quickly, but the
transfer as a whole
+ * never finishes on any reasonable timescale.
+ */
+ public static class SlowTrickleHttpHandler implements HttpHandler {
+
+ @Override
+ public void handle(HttpExchange exchange) throws IOException {
+ // 0 body length with no prior Content-Length header tells the JDK
HTTP server to use
+ // chunked transfer encoding, so the client never sees a declared
size to fail fast on.
+ exchange.sendResponseHeaders(HttpURLConnection.HTTP_OK, 0);
+ var body = exchange.getResponseBody();
+ try {
+ while (true) {
+ body.write(0);
+ body.flush();
+ Thread.sleep(50);
+ }
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ } finally {
+ exchange.close();
+ }
+ }
+ }
+
+ /**
+ * Handler that sends response headers declaring a body, then stalls
without sending any body
+ * bytes, so a client read blocks until the socket read timeout fires.
+ */
+ public static class StallAfterHeadersHttpHandler implements HttpHandler {
+
+ @Override
+ public void handle(HttpExchange exchange) throws IOException {
+ // Declare a fixed-length body but never write it; the client's
read blocks on the body.
+ exchange.sendResponseHeaders(HttpURLConnection.HTTP_OK, 1024);
+ try {
+ Thread.sleep(5_000);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ } finally {
+ exchange.close();
+ }
+ }
+ }
}