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]

Reply via email to