This is an automated email from the ASF dual-hosted git repository.
sarutak pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/spark.git
The following commit(s) were added to refs/heads/master by this push:
new 97535e4d70cb [SPARK-57892][CORE] Add TokenIngestor interface and
FileTokenIngestor implementation
97535e4d70cb is described below
commit 97535e4d70cbdc10487fda7aae4351587c215c36
Author: Shrirang Mhalgi <[email protected]>
AuthorDate: Fri Jul 17 09:08:52 2026 +0900
[SPARK-57892][CORE] Add TokenIngestor interface and FileTokenIngestor
implementation
### What changes were proposed in this pull request?
This PR adds `TokenIngestor` interface and `FileTokenIngestor`
implementation for OIDC credential propagation framework (SPARK-57703).
1. `TokenIngestor`: Added `DeveloperApi` interface containing single method
`load()` that returns `Optional<UserContext>`, providing the abstraction for
reading identity tokens.
2. `FileTokenIngestor`: `FileTokenIngestor` reads an OIDC JWT from a file
path (e.g. Kubernetes projected service account token at
`/var/run/secrets/tokens/spark-identity`), parses claims (sub, iss, iat, exp)
into a `UserContext`, detects file rotation via `mtime` caching, and handles
errors gracefully (missing file, empty content, malformed JWT all return
Optional.empty() without throwing).
### Why are the changes needed?
This is subtask 3 of the OIDC Credential Propagation SPIP (SPARK-57703).
The `UserCredentialManager` (subtask 4) needs a way to obtain the current
identity token from the local filesystem. On Kubernetes, identity tokens are
delivered as projected volumes that are periodically rotated by kubelet.
`FileTokenIngestor` will act as a bridge between these filesystem-delivered
tokens and the credential propagation framework, supporting both workload-level
service account tokens and per-user [...]
### Does this PR introduce _any_ user-facing change?
No. These are internal SPIs annotated DeveloperApi that will be consumed by
`UserCredentialManager` in a subsequent OIDC PRs. No configuration keys are
activated and no existing behavior is changed.
### How was this patch tested?
Added unit tests in `FileTokenIngestorSuite` covering:
1. Valid token parsing (claims extraction, UserContext construction)
2. File rotation detection (mtime-based re-parse)
3. Missing file handling
4. Empty file and whitespace-only file handling
5. Malformed JWT handling
6. Missing required claims (`sub`, `iss`)
7. Optional claims absent (iat, exp → null Instant)
8. Mtime caching (same file returns cached result)
9. Kubernetes projected SA token format
10. Per-user identity token format
### Was this patch authored or co-authored using generative AI tooling?
No.
Closes #57262 from shrirangmhalgi/SPARK-57892-token-ingestor.
Authored-by: Shrirang Mhalgi <[email protected]>
Signed-off-by: Kousuke Saruta <[email protected]>
---
.../apache/spark/security/FileTokenIngestor.java | 156 ++++++++++++
.../org/apache/spark/security/TokenIngestor.java | 43 ++++
.../spark/security/FileTokenIngestorSuite.java | 263 +++++++++++++++++++++
3 files changed, 462 insertions(+)
diff --git
a/core/src/main/java/org/apache/spark/security/FileTokenIngestor.java
b/core/src/main/java/org/apache/spark/security/FileTokenIngestor.java
new file mode 100644
index 000000000000..c3fd8d65a7a0
--- /dev/null
+++ b/core/src/main/java/org/apache/spark/security/FileTokenIngestor.java
@@ -0,0 +1,156 @@
+/*
+ * 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.spark.security;
+
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.time.Instant;
+import java.util.Base64;
+import java.util.Optional;
+
+import com.fasterxml.jackson.databind.JsonNode;
+import com.fasterxml.jackson.databind.ObjectMapper;
+
+import org.apache.spark.annotation.Private;
+import org.apache.spark.internal.LogKeys;
+import org.apache.spark.internal.MDC;
+import org.apache.spark.internal.SparkLogger;
+import org.apache.spark.internal.SparkLoggerFactory;
+
+/**
+ * A {@link TokenIngestor} that reads an OIDC identity token from a file.
+ * <p>
+ * The file path is typically a Kubernetes projected service account token
+ * (e.g., {@code /var/run/secrets/tokens/spark-identity}) or a path configured
via
+ * {@code spark.security.credentials.identityToken.file}.
+ * <p>
+ * This implementation:
+ * <ul>
+ * <li>Detects file rotation via mtime change (only re-parses when the file
changes)</li>
+ * <li>Parses JWT claims by Base64-decoding the payload segment directly,
without
+ * signature verification, since the token is trusted from the local
filesystem
+ * and works with both signed (RS256, ES256) and unsigned tokens</li>
+ * <li>Handles errors gracefully: malformed JWT, missing file, or empty
content
+ * returns empty rather than throwing</li>
+ * </ul>
+ *
+ * @since 4.3.0
+ */
+@Private
+public class FileTokenIngestor implements TokenIngestor {
+
+ private static final SparkLogger LOG =
SparkLoggerFactory.getLogger(FileTokenIngestor.class);
+ private static final ObjectMapper MAPPER = new ObjectMapper();
+
+ private final Path tokenPath;
+
+ // Cached state for rotation detection.
+ // Write order matters for thread-safety: cachedContext must be visible
before
+ // lastMtime, so a concurrent reader never sees a new mtime with a stale
context.
+ private volatile UserContext cachedContext = null;
+ private volatile long lastMtime = -1L;
+
+ /**
+ * Construct a new FileTokenIngestor.
+ *
+ * @param tokenPath path to the identity token file (must not be null)
+ */
+ public FileTokenIngestor(Path tokenPath) {
+ this.tokenPath = tokenPath;
+ }
+
+ @Override
+ public Optional<UserContext> load() {
+ try {
+ if (!Files.exists(tokenPath)) {
+ LOG.debug("Token file does not exist: {}", tokenPath);
+ return Optional.empty();
+ }
+
+ // File did not change since last successful parse
+ long currentMtime = Files.getLastModifiedTime(tokenPath).toMillis();
+ if (currentMtime == lastMtime && cachedContext != null) {
+ return Optional.of(cachedContext);
+ }
+
+ String content = new String(Files.readAllBytes(tokenPath),
StandardCharsets.UTF_8).trim();
+ if (content.isEmpty()) {
+ LOG.warn("Token file is empty: {}", MDC.of(LogKeys.PATH, tokenPath));
+ return Optional.empty();
+ }
+
+ Optional<UserContext> userContext = parseJwt(content);
+ if (userContext.isPresent()) {
+ // If the new file has invalid content, return empty and do NOT
+ // fall back to the previously cached context. The caller will retry
on next poll.
+ cachedContext = userContext.get();
+ lastMtime = currentMtime;
+ }
+ return userContext;
+ } catch (Exception e) {
+ LOG.warn("Failed to load token from {}: {}", e,
+ MDC.of(LogKeys.PATH, tokenPath),
+ MDC.of(LogKeys.CLASS_NAME, e.getClass().getName()));
+ return Optional.empty();
+ }
+ }
+
+ /**
+ * Parse a JWT token string into a UserContext by Base64-decoding the
payload segment.
+ * This works with both signed (RS256, ES256) and unsigned (alg:none) tokens
since
+ * we never verify the signature - the token is trusted from the local
filesystem and
+ * will be re-verified downstream at the STS token exchange.
+ */
+ private Optional<UserContext> parseJwt(String token) {
+ try {
+ String[] parts = token.split("\\.");
+ if (parts.length < 2) {
+ LOG.warn("JWT token does not have the expected header.payload format");
+ return Optional.empty();
+ }
+
+ // Decode the payload (second segment) - works for both signed and
unsigned JWTs
+ byte[] payloadBytes = Base64.getUrlDecoder().decode(parts[1]);
+ JsonNode claims = MAPPER.readTree(payloadBytes);
+
+ String subject = claims.has("sub") ? claims.get("sub").asText() : null;
+ String issuer = claims.has("iss") ? claims.get("iss").asText() : null;
+
+ if (subject == null || subject.isEmpty()) {
+ LOG.warn("JWT token missing required 'sub' claim");
+ return Optional.empty();
+ }
+ if (issuer == null || issuer.isEmpty()) {
+ LOG.warn("JWT token missing required 'iss' claim");
+ return Optional.empty();
+ }
+
+ Instant issuedAt = claims.has("iat")
+ ? Instant.ofEpochSecond(claims.get("iat").asLong()) : null;
+ Instant expiresAt = claims.has("exp")
+ ? Instant.ofEpochSecond(claims.get("exp").asLong()) : null;
+
+ return Optional.of(new UserContext(subject, issuer, token, issuedAt,
expiresAt));
+ } catch (Exception e) {
+ LOG.warn("Failed to parse JWT token: {}",
+ MDC.of(LogKeys.CLASS_NAME, e.getClass().getName()));
+ return Optional.empty();
+ }
+ }
+}
diff --git a/core/src/main/java/org/apache/spark/security/TokenIngestor.java
b/core/src/main/java/org/apache/spark/security/TokenIngestor.java
new file mode 100644
index 000000000000..ad4a208e663f
--- /dev/null
+++ b/core/src/main/java/org/apache/spark/security/TokenIngestor.java
@@ -0,0 +1,43 @@
+/*
+ * 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.spark.security;
+
+import java.util.Optional;
+
+import org.apache.spark.annotation.DeveloperApi;
+
+/**
+ * :: DeveloperApi ::
+ * Read an OIDC identity token and produces a {@link UserContext}.
+ * <p>
+ * Implementation should be stateless with respect to Spark configuration;
+ * configuration is passed at construction time.
+ *
+ * @since 4.3.0
+ */
+@DeveloperApi
+public interface TokenIngestor {
+
+ /**
+ * Attempt to load the current identity token and parse it into a
UserContext.
+ *
+ * @return a present Optional containing the UserContext if a valid token is
available,
+ * or empty if unavailable (e.g. empty content / missing file).
+ */
+ Optional<UserContext> load();
+}
diff --git
a/core/src/test/java/org/apache/spark/security/FileTokenIngestorSuite.java
b/core/src/test/java/org/apache/spark/security/FileTokenIngestorSuite.java
new file mode 100644
index 000000000000..f6f8ea2c3636
--- /dev/null
+++ b/core/src/test/java/org/apache/spark/security/FileTokenIngestorSuite.java
@@ -0,0 +1,263 @@
+/*
+ * 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.spark.security;
+
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.security.KeyPair;
+import java.security.KeyPairGenerator;
+import java.util.Comparator;
+import java.util.Date;
+import java.util.Optional;
+
+import io.jsonwebtoken.Jwts;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.*;
+
+public class FileTokenIngestorSuite {
+
+ private Path tempDir;
+
+ @BeforeEach
+ public void setUp() throws Exception {
+ tempDir = Files.createTempDirectory("token-ingestor-test");
+ }
+
+ @AfterEach
+ public void tearDown() throws Exception {
+ if (tempDir != null) {
+ Files.walk(tempDir).sorted(Comparator.reverseOrder())
+ .forEach(path -> {
+ try { Files.deleteIfExists(path); } catch (Exception e) { /*
cleanup, ignore */ }
+ });
+ }
+ }
+
+ private String createUnsignedJwt(String subject, String issuer, Date
issuedAt, Date expiresAt) {
+ var builder = Jwts.builder()
+ .subject(subject)
+ .issuer(issuer);
+ if (issuedAt != null) builder.issuedAt(issuedAt);
+ if (expiresAt != null) builder.expiration(expiresAt);
+ return builder.compact();
+ }
+
+ private String createUnsignedJwt(String subject, String issuer) {
+ return createUnsignedJwt(subject, issuer,
+ new Date(), new Date(System.currentTimeMillis() + 3600000));
+ }
+
+ private String createUnsignedJwt() {
+ return createUnsignedJwt("system:serviceaccount:ns:sa",
+ "https://kubernetes.default.svc");
+ }
+
+ private String createSignedJwt(String subject, String issuer) throws
Exception {
+ KeyPairGenerator kpg = KeyPairGenerator.getInstance("RSA");
+ kpg.initialize(2048);
+ KeyPair kp = kpg.generateKeyPair();
+ return Jwts.builder()
+ .subject(subject)
+ .issuer(issuer)
+ .issuedAt(new Date())
+ .expiration(new Date(System.currentTimeMillis() + 3600000))
+ .signWith(kp.getPrivate())
+ .compact();
+ }
+
+ private void writeToken(Path path, String content) throws Exception {
+ Files.write(path, content.getBytes("UTF-8"));
+ }
+
+ @Test
+ public void loadReturnsValidUserContextFromTokenFile() throws Exception {
+ Path tokenFile = tempDir.resolve("token");
+ String token = createUnsignedJwt();
+ writeToken(tokenFile, token);
+
+ FileTokenIngestor ingestor = new FileTokenIngestor(tokenFile);
+ Optional<UserContext> result = ingestor.load();
+
+ assertTrue(result.isPresent());
+ UserContext ctx = result.get();
+ assertEquals("system:serviceaccount:ns:sa", ctx.getPrincipal());
+ assertEquals("https://kubernetes.default.svc", ctx.getIssuer());
+ assertEquals(token, ctx.getRawToken());
+ assertNotNull(ctx.getIssuedAt());
+ assertNotNull(ctx.getExpiresAt());
+ }
+
+ @Test
+ public void loadWorksWithSignedJwt() throws Exception {
+ Path tokenFile = tempDir.resolve("token");
+ String token = createSignedJwt("[email protected]",
"https://oidc.eks.us-west-2.amazonaws.com/id/ABC123");
+ writeToken(tokenFile, token);
+
+ FileTokenIngestor ingestor = new FileTokenIngestor(tokenFile);
+ Optional<UserContext> result = ingestor.load();
+
+ assertTrue(result.isPresent());
+ UserContext ctx = result.get();
+ assertEquals("[email protected]", ctx.getPrincipal());
+ assertEquals("https://oidc.eks.us-west-2.amazonaws.com/id/ABC123",
ctx.getIssuer());
+ assertEquals(token, ctx.getRawToken());
+ assertNotNull(ctx.getIssuedAt());
+ assertNotNull(ctx.getExpiresAt());
+ }
+
+ @Test
+ public void loadDetectsFileRotation() throws Exception {
+ Path tokenFile = tempDir.resolve("token");
+ String token1 = createUnsignedJwt("user1", "https://issuer.example.com");
+ writeToken(tokenFile, token1);
+
+ FileTokenIngestor ingestor = new FileTokenIngestor(tokenFile);
+ Optional<UserContext> result1 = ingestor.load();
+ assertEquals("user1", result1.get().getPrincipal());
+
+ // Simulate token rotation: sleep to ensure mtime differs
+ Thread.sleep(100);
+ String token2 = createUnsignedJwt("user2", "https://issuer.example.com");
+ writeToken(tokenFile, token2);
+
+ Optional<UserContext> result2 = ingestor.load();
+ assertEquals("user2", result2.get().getPrincipal());
+ assertEquals(token2, result2.get().getRawToken());
+ }
+
+ @Test
+ public void loadReturnsEmptyForMissingFile() {
+ Path tokenFile = tempDir.resolve("nonexistent");
+ FileTokenIngestor ingestor = new FileTokenIngestor(tokenFile);
+ assertTrue(ingestor.load().isEmpty());
+ }
+
+ @Test
+ public void loadReturnsEmptyForEmptyFile() throws Exception {
+ Path tokenFile = tempDir.resolve("token");
+ writeToken(tokenFile, "");
+
+ FileTokenIngestor ingestor = new FileTokenIngestor(tokenFile);
+ assertTrue(ingestor.load().isEmpty());
+ }
+
+ @Test
+ public void loadReturnsEmptyForWhitespaceOnlyFile() throws Exception {
+ Path tokenFile = tempDir.resolve("token");
+ writeToken(tokenFile, " \n\t ");
+
+ FileTokenIngestor ingestor = new FileTokenIngestor(tokenFile);
+ assertTrue(ingestor.load().isEmpty());
+ }
+
+ @Test
+ public void loadReturnsEmptyForMalformedJwt() throws Exception {
+ Path tokenFile = tempDir.resolve("token");
+ writeToken(tokenFile, "not-a-jwt");
+
+ FileTokenIngestor ingestor = new FileTokenIngestor(tokenFile);
+ assertTrue(ingestor.load().isEmpty());
+ }
+
+ @Test
+ public void loadReturnsEmptyForJwtMissingSubClaim() throws Exception {
+ Path tokenFile = tempDir.resolve("token");
+ String token = Jwts.builder()
+ .issuer("https://issuer.example.com")
+ .compact();
+ writeToken(tokenFile, token);
+
+ FileTokenIngestor ingestor = new FileTokenIngestor(tokenFile);
+ assertTrue(ingestor.load().isEmpty());
+ }
+
+ @Test
+ public void loadReturnsEmptyForJwtMissingIssClaim() throws Exception {
+ Path tokenFile = tempDir.resolve("token");
+ String token = Jwts.builder()
+ .subject("user1")
+ .compact();
+ writeToken(tokenFile, token);
+
+ FileTokenIngestor ingestor = new FileTokenIngestor(tokenFile);
+ assertTrue(ingestor.load().isEmpty());
+ }
+
+ @Test
+ public void loadHandlesJwtWithoutIatAndExpClaims() throws Exception {
+ Path tokenFile = tempDir.resolve("token");
+ String token = Jwts.builder()
+ .subject("user1")
+ .issuer("https://issuer.example.com")
+ .compact();
+ writeToken(tokenFile, token);
+
+ FileTokenIngestor ingestor = new FileTokenIngestor(tokenFile);
+ Optional<UserContext> result = ingestor.load();
+ assertTrue(result.isPresent());
+ assertEquals("user1", result.get().getPrincipal());
+ assertNull(result.get().getIssuedAt());
+ assertNull(result.get().getExpiresAt());
+ }
+
+ @Test
+ public void loadUsesCachedResultWhenFileUnchanged() throws Exception {
+ Path tokenFile = tempDir.resolve("token");
+ String token = createUnsignedJwt("cached-user",
"https://issuer.example.com");
+ writeToken(tokenFile, token);
+
+ FileTokenIngestor ingestor = new FileTokenIngestor(tokenFile);
+ Optional<UserContext> result1 = ingestor.load();
+ Optional<UserContext> result2 = ingestor.load();
+
+ // Reference identity proves the cached path was hit (not re-parsed)
+ assertSame(result1.get(), result2.get());
+ assertEquals("cached-user", result1.get().getPrincipal());
+ }
+
+ @Test
+ public void loadWorksWithKubernetesProjectedToken() throws Exception {
+ Path tokenFile = tempDir.resolve("token");
+ String token = createUnsignedJwt(
+ "system:serviceaccount:spark-ns:spark-driver-sa",
+
"https://oidc.eks.us-west-2.amazonaws.com/id/EXAMPLED539D4633E53DE1B716D3041E");
+ writeToken(tokenFile, token);
+
+ FileTokenIngestor ingestor = new FileTokenIngestor(tokenFile);
+ Optional<UserContext> result = ingestor.load();
+ assertTrue(result.isPresent());
+ assertEquals("system:serviceaccount:spark-ns:spark-driver-sa",
result.get().getPrincipal());
+ assertTrue(result.get().getIssuer().contains("oidc.eks"));
+ }
+
+ @Test
+ public void loadWorksWithPerUserIdentityToken() throws Exception {
+ Path tokenFile = tempDir.resolve("token");
+ String token = createUnsignedJwt("[email protected]",
"https://accounts.google.com");
+ writeToken(tokenFile, token);
+
+ FileTokenIngestor ingestor = new FileTokenIngestor(tokenFile);
+ Optional<UserContext> result = ingestor.load();
+ assertTrue(result.isPresent());
+ assertEquals("[email protected]", result.get().getPrincipal());
+ assertEquals("https://accounts.google.com", result.get().getIssuer());
+ }
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]