This is an automated email from the ASF dual-hosted git repository.
rzo1 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/storm.git
The following commit(s) were added to refs/heads/master by this push:
new 96f65ace1 Migrate to Java 24+ compatible security APIs and add Java
25 to CI (#8483)
96f65ace1 is described below
commit 96f65ace17487ce1d0960553b0e526c105afa067
Author: Richard Zowalla <[email protected]>
AuthorDate: Thu Apr 2 09:42:25 2026 +0200
Migrate to Java 24+ compatible security APIs and add Java 25 to CI
(#8483)
* Migrate to Java 24+ compatible security APIs and add Java 25 to CI (#8456)
- Replace Subject.doAs(), Subject.getSubject(AccessControlContext),
AccessController.getContext(), and System.getSecurityManager() — all
removed in Java 24 with runtime-dispatched compatibility shims that
work on Java 17 through 25+.
- Introduce SubjectCompat utility with MethodHandle-based dispatch:
currentSubject() maps to Subject.current() (18+) or
Subject.getSubject(AccessController.getContext()) (17);
doAs() maps to Subject.callAs() (18+) or Subject.doAs() (17).
- Migrate all 12 removed-API call sites across 8 files.
- Consolidate existing ReqContext MethodHandle shim into SubjectCompat.
- Remove dead SecurityManager code in NettyRenameThreadFactory.
- Add Java 25 to the GitHub Actions CI matrix.
* Checkstyle
---
.github/workflows/maven.yaml | 2 +-
.../jvm/org/apache/storm/daemon/worker/Worker.java | 4 +-
.../messaging/netty/KerberosSaslNettyClient.java | 60 +++----
.../messaging/netty/KerberosSaslNettyServer.java | 62 +++----
.../messaging/netty/NettyRenameThreadFactory.java | 4 +-
.../org/apache/storm/security/auth/ReqContext.java | 63 +------
.../apache/storm/security/auth/SubjectCompat.java | 183 +++++++++++++++++++++
.../auth/kerberos/KerberosSaslTransportPlugin.java | 60 +++----
.../WorkerTokenClientCallbackHandler.java | 4 +-
.../org/apache/storm/utils/HadoopLoginUtil.java | 4 +-
.../storm/security/auth/SubjectCompatTest.java | 70 ++++++++
.../org/apache/storm/security/auth/AuthTest.java | 8 +-
12 files changed, 346 insertions(+), 178 deletions(-)
diff --git a/.github/workflows/maven.yaml b/.github/workflows/maven.yaml
index f95b502ae..671d06139 100644
--- a/.github/workflows/maven.yaml
+++ b/.github/workflows/maven.yaml
@@ -28,7 +28,7 @@ jobs:
strategy:
matrix:
os: [ ubuntu-latest ]
- java: [ 17, 21 ]
+ java: [ 17, 21, 25 ]
module: [ Client, Server, Core, External, Integration-Test ]
experimental: [false]
fail-fast: false
diff --git a/storm-client/src/jvm/org/apache/storm/daemon/worker/Worker.java
b/storm-client/src/jvm/org/apache/storm/daemon/worker/Worker.java
index 988e7bd72..e78b01804 100644
--- a/storm-client/src/jvm/org/apache/storm/daemon/worker/Worker.java
+++ b/storm-client/src/jvm/org/apache/storm/daemon/worker/Worker.java
@@ -20,7 +20,6 @@ import java.io.File;
import java.io.IOException;
import java.net.UnknownHostException;
import java.nio.charset.Charset;
-import java.security.PrivilegedExceptionAction;
import java.util.ArrayList;
import java.util.Collection;
import java.util.HashMap;
@@ -58,6 +57,7 @@ import org.apache.storm.messaging.IContext;
import org.apache.storm.metrics2.StormMetricRegistry;
import org.apache.storm.security.auth.ClientAuthUtils;
import org.apache.storm.security.auth.IAutoCredentials;
+import org.apache.storm.security.auth.SubjectCompat;
import org.apache.storm.shade.com.google.common.base.Preconditions;
import org.apache.storm.shade.org.apache.commons.io.FileUtils;
import org.apache.storm.shade.org.apache.commons.lang3.ObjectUtils;
@@ -202,7 +202,7 @@ public class Worker implements Shutdownable, DaemonCommon {
autoCreds = ClientAuthUtils.getAutoCredentials(topologyConf);
subject = ClientAuthUtils.populateSubject(null, autoCreds, initCreds);
- Subject.doAs(subject, (PrivilegedExceptionAction<Object>)
+ SubjectCompat.doAs(subject,
() -> loadWorker(stateStorage, stormClusterState, initCreds,
initialCredentials)
);
diff --git
a/storm-client/src/jvm/org/apache/storm/messaging/netty/KerberosSaslNettyClient.java
b/storm-client/src/jvm/org/apache/storm/messaging/netty/KerberosSaslNettyClient.java
index 31425d4b4..6fec202ed 100644
---
a/storm-client/src/jvm/org/apache/storm/messaging/netty/KerberosSaslNettyClient.java
+++
b/storm-client/src/jvm/org/apache/storm/messaging/netty/KerberosSaslNettyClient.java
@@ -14,8 +14,6 @@ package org.apache.storm.messaging.netty;
import java.io.IOException;
import java.security.Principal;
-import java.security.PrivilegedActionException;
-import java.security.PrivilegedExceptionAction;
import java.util.Map;
import java.util.TreeMap;
import javax.security.auth.Subject;
@@ -28,6 +26,7 @@ import javax.security.sasl.Sasl;
import javax.security.sasl.SaslClient;
import javax.security.sasl.SaslException;
import org.apache.storm.security.auth.ClientAuthUtils;
+import org.apache.storm.security.auth.SubjectCompat;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -95,28 +94,25 @@ public class KerberosSaslNettyClient {
final String fServiceName = serviceName;
final CallbackHandler fch = ch;
LOG.debug("Kerberos Client with principal: {}, host: {}",
fPrincipalName, fHost);
- saslClient = Subject.doAs(subject, new
PrivilegedExceptionAction<SaslClient>() {
- @Override
- public SaslClient run() {
- try {
- Map<String, String> props = new TreeMap<String,
String>();
- props.put(Sasl.QOP, "auth");
- props.put(Sasl.SERVER_AUTH, "false");
- return Sasl.createSaslClient(
- new String[]{ SaslUtils.KERBEROS },
- fPrincipalName,
- fServiceName,
- fHost,
- props, fch);
- } catch (Exception e) {
- LOG.error("Subject failed to create sasl client.", e);
- return null;
- }
+ saslClient = SubjectCompat.doAs(subject, () -> {
+ try {
+ Map<String, String> props = new TreeMap<String, String>();
+ props.put(Sasl.QOP, "auth");
+ props.put(Sasl.SERVER_AUTH, "false");
+ return Sasl.createSaslClient(
+ new String[]{ SaslUtils.KERBEROS },
+ fPrincipalName,
+ fServiceName,
+ fHost,
+ props, fch);
+ } catch (Exception e) {
+ LOG.error("Subject failed to create sasl client.", e);
+ return null;
}
});
LOG.info("Got Client: {}", saslClient);
- } catch (PrivilegedActionException e) {
+ } catch (Exception e) {
LOG.error("KerberosSaslNettyClient: Could not create Sasl Netty
Client.");
throw new RuntimeException(e);
}
@@ -134,23 +130,17 @@ public class KerberosSaslNettyClient {
*/
public byte[] saslResponse(SaslMessageToken saslTokenMessage) {
try {
- final SaslMessageToken fSaslTokenMessage = saslTokenMessage;
- byte[] retval = Subject.doAs(subject, new
PrivilegedExceptionAction<byte[]>() {
- @Override
- public byte[] run() {
- try {
- byte[] retval =
saslClient.evaluateChallenge(fSaslTokenMessage
-
.getSaslToken());
- return retval;
- } catch (SaslException e) {
- LOG.error("saslResponse: Failed to respond to SASL
server's token:",
- e);
- throw new RuntimeException(e);
- }
+ return SubjectCompat.doAs(subject, () -> {
+ try {
+ return
saslClient.evaluateChallenge(saslTokenMessage.getSaslToken());
+ } catch (SaslException e) {
+ LOG.error("saslResponse: Failed to respond to SASL
server's token:", e);
+ throw new RuntimeException(e);
}
});
- return retval;
- } catch (PrivilegedActionException e) {
+ } catch (RuntimeException e) {
+ throw e;
+ } catch (Exception e) {
LOG.error("Failed to generate response for token: ", e);
throw new RuntimeException(e);
}
diff --git
a/storm-client/src/jvm/org/apache/storm/messaging/netty/KerberosSaslNettyServer.java
b/storm-client/src/jvm/org/apache/storm/messaging/netty/KerberosSaslNettyServer.java
index fefcdc66f..d3652e25e 100644
---
a/storm-client/src/jvm/org/apache/storm/messaging/netty/KerberosSaslNettyServer.java
+++
b/storm-client/src/jvm/org/apache/storm/messaging/netty/KerberosSaslNettyServer.java
@@ -14,8 +14,6 @@ package org.apache.storm.messaging.netty;
import java.io.IOException;
import java.security.Principal;
-import java.security.PrivilegedActionException;
-import java.security.PrivilegedExceptionAction;
import java.util.List;
import java.util.Map;
import java.util.TreeMap;
@@ -32,6 +30,7 @@ import javax.security.sasl.SaslException;
import javax.security.sasl.SaslServer;
import org.apache.storm.security.auth.ClientAuthUtils;
import org.apache.storm.security.auth.KerberosPrincipalToLocal;
+import org.apache.storm.security.auth.SubjectCompat;
import org.apache.storm.shade.org.apache.zookeeper.server.auth.KerberosName;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -83,26 +82,22 @@ class KerberosSaslNettyServer {
final String hostName = kerberosName.getHostName();
final String serviceName = kerberosName.getServiceName();
LOG.debug("Server with host: {}", hostName);
- saslServer =
- Subject.doAs(subject, new
PrivilegedExceptionAction<SaslServer>() {
- @Override
- public SaslServer run() {
- try {
- Map<String, String> props = new TreeMap<String,
String>();
- props.put(Sasl.QOP, "auth");
- props.put(Sasl.SERVER_AUTH, "false");
- return Sasl.createSaslServer(SaslUtils.KERBEROS,
- serviceName,
- hostName, props, fch);
- } catch (Exception e) {
- LOG.error("Subject failed to create sasl server.",
e);
- return null;
- }
- }
- });
+ saslServer = SubjectCompat.doAs(subject, () -> {
+ try {
+ Map<String, String> props = new TreeMap<String, String>();
+ props.put(Sasl.QOP, "auth");
+ props.put(Sasl.SERVER_AUTH, "false");
+ return Sasl.createSaslServer(SaslUtils.KERBEROS,
+ serviceName,
+ hostName, props, fch);
+ } catch (Exception e) {
+ LOG.error("Subject failed to create sasl server.", e);
+ return null;
+ }
+ });
LOG.info("Got Server: {}", saslServer);
- } catch (PrivilegedActionException e) {
+ } catch (Exception e) {
LOG.error("KerberosSaslNettyServer: Could not create SaslServer:
", e);
throw new RuntimeException(e);
}
@@ -124,23 +119,20 @@ class KerberosSaslNettyServer {
*/
public byte[] response(final byte[] token) {
try {
- byte[] retval = Subject.doAs(subject, new
PrivilegedExceptionAction<byte[]>() {
- @Override
- public byte[] run() {
- try {
- LOG.debug("response: Responding to input token of
length: {}",
- token.length);
- byte[] retval = saslServer.evaluateResponse(token);
- return retval;
- } catch (SaslException e) {
- LOG.error("response: Failed to evaluate client token
of length: {} : {}",
- token.length, e);
- throw new RuntimeException(e);
- }
+ return SubjectCompat.doAs(subject, () -> {
+ try {
+ LOG.debug("response: Responding to input token of length:
{}",
+ token.length);
+ return saslServer.evaluateResponse(token);
+ } catch (SaslException e) {
+ LOG.error("response: Failed to evaluate client token of
length: {} : {}",
+ token.length, e);
+ throw new RuntimeException(e);
}
});
- return retval;
- } catch (PrivilegedActionException e) {
+ } catch (RuntimeException e) {
+ throw e;
+ } catch (Exception e) {
LOG.error("Failed to generate response for token: ", e);
throw new RuntimeException(e);
}
diff --git
a/storm-client/src/jvm/org/apache/storm/messaging/netty/NettyRenameThreadFactory.java
b/storm-client/src/jvm/org/apache/storm/messaging/netty/NettyRenameThreadFactory.java
index f66510e70..1bc59c83a 100644
---
a/storm-client/src/jvm/org/apache/storm/messaging/netty/NettyRenameThreadFactory.java
+++
b/storm-client/src/jvm/org/apache/storm/messaging/netty/NettyRenameThreadFactory.java
@@ -24,9 +24,7 @@ public class NettyRenameThreadFactory implements
ThreadFactory {
private final String name;
public NettyRenameThreadFactory(String name) {
- SecurityManager s = System.getSecurityManager();
- group = (s != null) ? s.getThreadGroup() :
- Thread.currentThread().getThreadGroup();
+ group = Thread.currentThread().getThreadGroup();
this.name = name;
}
diff --git
a/storm-client/src/jvm/org/apache/storm/security/auth/ReqContext.java
b/storm-client/src/jvm/org/apache/storm/security/auth/ReqContext.java
index 830165084..3e61daf44 100644
--- a/storm-client/src/jvm/org/apache/storm/security/auth/ReqContext.java
+++ b/storm-client/src/jvm/org/apache/storm/security/auth/ReqContext.java
@@ -18,9 +18,6 @@
package org.apache.storm.security.auth;
-import java.lang.invoke.MethodHandle;
-import java.lang.invoke.MethodHandles;
-import java.lang.invoke.MethodType;
import java.net.InetAddress;
import java.security.Principal;
import java.util.Set;
@@ -38,8 +35,6 @@ import
org.apache.storm.shade.com.google.common.annotations.VisibleForTesting;
*/
public class ReqContext {
- private static final MethodHandle CURRENT = lookupCurrent();
-
private static final AtomicInteger uniqueId = new AtomicInteger(0);
//each thread will have its own request context
private static final ThreadLocal<ReqContext> ctxt =
@@ -167,64 +162,12 @@ public class ReqContext {
}
/**
- * Maps to Subject.current() is available, otherwise maps to
Subject.getSubject()
+ * Maps to Subject.current() if available, otherwise maps to
Subject.getSubject().
* @return the current subject
+ * @see SubjectCompat#currentSubject()
*/
public static Subject currentSubject() {
- try {
- return (Subject) CURRENT.invoke();
- } catch (Throwable t) {
- throw new RuntimeException(t);
- }
- }
-
- private static MethodHandle lookupCurrent() {
- final MethodHandles.Lookup lookup = MethodHandles.lookup();
- try {
- // Subject.getSubject(AccessControlContext) is deprecated for
removal and replaced by
- // Subject.current().
- // Lookup first the new API, since for Java versions where both
exists, the
- // new API delegates to the old API (for example Java 18, 19 and
20).
- // Otherwise (Java 17), lookup the old API.
- return lookup.findStatic(Subject.class, "current",
- MethodType.methodType(Subject.class));
- } catch (NoSuchMethodException e) {
- final MethodHandle getContext = lookupGetContext();
- final MethodHandle getSubject = lookupGetSubject();
- return MethodHandles.filterReturnValue(getContext, getSubject);
- } catch (IllegalAccessException e) {
- throw new AssertionError(e);
- }
- }
-
- private static MethodHandle lookupGetSubject() {
- final MethodHandles.Lookup lookup = MethodHandles.lookup();
- try {
- final Class<?> contextClazz =
- ClassLoader.getSystemClassLoader()
- .loadClass("java.security.AccessControlContext");
- return lookup.findStatic(Subject.class, "getSubject",
- MethodType.methodType(Subject.class, contextClazz));
- } catch (ClassNotFoundException | NoSuchMethodException |
IllegalAccessException e) {
- throw new AssertionError(e);
- }
- }
-
- private static MethodHandle lookupGetContext() {
- try {
- // Use reflection to work with Java versions that have and don't
have AccessController.
- final Class<?> controllerClazz =
-
ClassLoader.getSystemClassLoader().loadClass("java.security.AccessController");
- final Class<?> contextClazz =
- ClassLoader.getSystemClassLoader()
- .loadClass("java.security.AccessControlContext");
-
- MethodHandles.Lookup lookup = MethodHandles.lookup();
- return lookup.findStatic(controllerClazz, "getContext",
- MethodType.methodType(contextClazz));
- } catch (ClassNotFoundException | NoSuchMethodException |
IllegalAccessException e) {
- throw new AssertionError(e);
- }
+ return SubjectCompat.currentSubject();
}
}
diff --git
a/storm-client/src/jvm/org/apache/storm/security/auth/SubjectCompat.java
b/storm-client/src/jvm/org/apache/storm/security/auth/SubjectCompat.java
new file mode 100644
index 000000000..a54541373
--- /dev/null
+++ b/storm-client/src/jvm/org/apache/storm/security/auth/SubjectCompat.java
@@ -0,0 +1,183 @@
+/*
+ * 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.storm.security.auth;
+
+import java.lang.invoke.MethodHandle;
+import java.lang.invoke.MethodHandles;
+import java.lang.invoke.MethodType;
+import java.util.concurrent.Callable;
+import java.util.concurrent.CompletionException;
+import javax.security.auth.Subject;
+
+/**
+ * Compatibility shim for {@link Subject} methods removed in Java 24.
+ *
+ * <ul>
+ * <li>{@code Subject.getSubject(AccessControlContext)} → {@code
Subject.current()} (Java 18+)</li>
+ * <li>{@code Subject.doAs(Subject, PrivilegedExceptionAction)} →
{@code Subject.callAs(Subject, Callable)} (Java 18+)</li>
+ * </ul>
+ *
+ * <p>All dispatch is resolved once at class-init via {@link MethodHandle}, so
there is no per-call reflection overhead.
+ */
+public final class SubjectCompat {
+
+ private static final MethodHandle CURRENT = lookupCurrent();
+ private static final MethodHandle CALL_AS = lookupCallAs();
+ private static final boolean USE_CALL_AS = probeCallAs();
+
+ private SubjectCompat() {
+ }
+
+ /**
+ * Return the current {@link Subject}, equivalent to {@code
Subject.current()} on Java 18+
+ * or {@code Subject.getSubject(AccessController.getContext())} on Java 17.
+ */
+ public static Subject currentSubject() {
+ try {
+ return (Subject) CURRENT.invoke();
+ } catch (Throwable t) {
+ throw new RuntimeException(t);
+ }
+ }
+
+ /**
+ * Execute a {@link Callable} as the given {@link Subject}, equivalent to
+ * {@code Subject.callAs(subject, callable)} on Java 18+ or
+ * {@code Subject.doAs(subject, privilegedAction)} on Java 17.
+ *
+ * @param subject the subject to run as (may be {@code null})
+ * @param callable the action to execute
+ * @param <T> return type
+ * @return the callable's result
+ * @throws Exception if the callable throws
+ */
+ @SuppressWarnings("unchecked")
+ public static <T> T doAs(Subject subject, Callable<T> callable) throws
Exception {
+ try {
+ if (USE_CALL_AS) {
+ return (T) CALL_AS.invoke(subject, callable);
+ } else {
+ // Java 17 path: wrap Callable in PrivilegedExceptionAction
for Subject.doAs()
+ return (T) CALL_AS.invoke(subject,
+ (java.security.PrivilegedExceptionAction<T>)
callable::call);
+ }
+ } catch (CompletionException e) {
+ // Subject.callAs wraps in CompletionException
+ unwrapAndThrow(e.getCause());
+ throw e; // unreachable
+ } catch (java.security.PrivilegedActionException e) {
+ // Subject.doAs wraps in PrivilegedActionException
+ unwrapAndThrow(e.getCause());
+ throw e; // unreachable
+ } catch (Exception e) {
+ throw e;
+ } catch (Throwable t) {
+ throw new RuntimeException(t);
+ }
+ }
+
+ private static void unwrapAndThrow(Throwable cause) throws Exception {
+ if (cause instanceof Exception) {
+ throw (Exception) cause;
+ }
+ if (cause instanceof Error) {
+ throw (Error) cause;
+ }
+ throw new RuntimeException(cause);
+ }
+
+ // --- MethodHandle lookups resolved once at class init ---
+
+ private static boolean probeCallAs() {
+ try {
+ MethodHandles.lookup().findStatic(Subject.class, "callAs",
+ MethodType.methodType(Object.class, Subject.class,
Callable.class));
+ return true;
+ } catch (NoSuchMethodException | IllegalAccessException e) {
+ return false;
+ }
+ }
+
+ private static MethodHandle lookupCallAs() {
+ MethodHandles.Lookup lookup = MethodHandles.lookup();
+ try {
+ // Java 18+: Subject.callAs(Subject, Callable)
+ return lookup.findStatic(Subject.class, "callAs",
+ MethodType.methodType(Object.class, Subject.class,
Callable.class));
+ } catch (NoSuchMethodException e) {
+ // Java 17: Subject.doAs(Subject, PrivilegedExceptionAction)
+ try {
+ return lookup.findStatic(Subject.class, "doAs",
+ MethodType.methodType(Object.class, Subject.class,
+ java.security.PrivilegedExceptionAction.class));
+ } catch (NoSuchMethodException | IllegalAccessException ex) {
+ throw new AssertionError(ex);
+ }
+ } catch (IllegalAccessException e) {
+ throw new AssertionError(e);
+ }
+ }
+
+ private static MethodHandle lookupCurrent() {
+ final MethodHandles.Lookup lookup = MethodHandles.lookup();
+ try {
+ // Subject.getSubject(AccessControlContext) is deprecated for
removal and replaced by
+ // Subject.current().
+ // Lookup first the new API, since for Java versions where both
exist, the
+ // new API delegates to the old API (for example Java 18, 19 and
20).
+ // Otherwise (Java 17), lookup the old API.
+ return lookup.findStatic(Subject.class, "current",
+ MethodType.methodType(Subject.class));
+ } catch (NoSuchMethodException e) {
+ final MethodHandle getContext = lookupGetContext();
+ final MethodHandle getSubject = lookupGetSubject();
+ return MethodHandles.filterReturnValue(getContext, getSubject);
+ } catch (IllegalAccessException e) {
+ throw new AssertionError(e);
+ }
+ }
+
+ private static MethodHandle lookupGetSubject() {
+ final MethodHandles.Lookup lookup = MethodHandles.lookup();
+ try {
+ final Class<?> contextClazz =
+ ClassLoader.getSystemClassLoader()
+ .loadClass("java.security.AccessControlContext");
+ return lookup.findStatic(Subject.class, "getSubject",
+ MethodType.methodType(Subject.class, contextClazz));
+ } catch (ClassNotFoundException | NoSuchMethodException |
IllegalAccessException e) {
+ throw new AssertionError(e);
+ }
+ }
+
+ private static MethodHandle lookupGetContext() {
+ try {
+ final Class<?> controllerClazz =
+
ClassLoader.getSystemClassLoader().loadClass("java.security.AccessController");
+ final Class<?> contextClazz =
+ ClassLoader.getSystemClassLoader()
+ .loadClass("java.security.AccessControlContext");
+ MethodHandles.Lookup lookup = MethodHandles.lookup();
+ return lookup.findStatic(controllerClazz, "getContext",
+ MethodType.methodType(contextClazz));
+ } catch (ClassNotFoundException | NoSuchMethodException |
IllegalAccessException e) {
+ throw new AssertionError(e);
+ }
+ }
+}
diff --git
a/storm-client/src/jvm/org/apache/storm/security/auth/kerberos/KerberosSaslTransportPlugin.java
b/storm-client/src/jvm/org/apache/storm/security/auth/kerberos/KerberosSaslTransportPlugin.java
index 980ac4478..395632363 100644
---
a/storm-client/src/jvm/org/apache/storm/security/auth/kerberos/KerberosSaslTransportPlugin.java
+++
b/storm-client/src/jvm/org/apache/storm/security/auth/kerberos/KerberosSaslTransportPlugin.java
@@ -14,8 +14,6 @@ package org.apache.storm.security.auth.kerberos;
import java.io.IOException;
import java.security.Principal;
-import java.security.PrivilegedActionException;
-import java.security.PrivilegedExceptionAction;
import java.util.Map;
import java.util.Set;
import java.util.SortedMap;
@@ -29,6 +27,7 @@ import javax.security.sasl.Sasl;
import org.apache.storm.generated.WorkerToken;
import org.apache.storm.messaging.netty.Login;
import org.apache.storm.security.auth.ClientAuthUtils;
+import org.apache.storm.security.auth.SubjectCompat;
import org.apache.storm.security.auth.sasl.SaslTransportPlugin;
import org.apache.storm.security.auth.sasl.SimpleSaslServerCallbackHandler;
import org.apache.storm.security.auth.workertoken.WorkerTokenAuthorizer;
@@ -200,23 +199,19 @@ public class KerberosSaslTransportPlugin extends
SaslTransportPlugin {
//open Sasl transport with the login credential
try {
- Subject.doAs(subject,
- new PrivilegedExceptionAction<Void>() {
- @Override
- public Void run() {
- try {
- LOG.debug("do as:" + principal);
- sasalTransport.open();
- } catch (Exception e) {
- LOG.error("Client failed to open
SaslClientTransport to interact with a server during "
- + "session initiation: "
- + e,
- e);
- }
- return null;
- }
- });
- } catch (PrivilegedActionException e) {
+ SubjectCompat.doAs(subject, () -> {
+ try {
+ LOG.debug("do as:" + principal);
+ sasalTransport.open();
+ } catch (Exception e) {
+ LOG.error("Client failed to open SaslClientTransport to
interact with a server during "
+ + "session initiation: "
+ + e,
+ e);
+ }
+ return null;
+ });
+ } catch (Exception e) {
throw new RuntimeException(e);
}
@@ -265,20 +260,19 @@ public class KerberosSaslTransportPlugin extends
SaslTransportPlugin {
@Override
public TTransport getTransport(final TTransport trans) {
try {
- return Subject.doAs(subject,
- (PrivilegedExceptionAction<TTransport>) () -> {
- try {
- return wrapped.getTransport(trans);
- } catch (Exception e) {
- LOG.debug("Storm server failed to open transport
to interact with a client during "
- + "session initiation: "
- + e,
- e);
- return new NoOpTTrasport(null);
- }
- });
- } catch (PrivilegedActionException e) {
- LOG.error("Storm server experienced a
PrivilegedActionException exception while creating a transport "
+ return SubjectCompat.doAs(subject, () -> {
+ try {
+ return wrapped.getTransport(trans);
+ } catch (Exception e) {
+ LOG.debug("Storm server failed to open transport to
interact with a client during "
+ + "session initiation: "
+ + e,
+ e);
+ return new NoOpTTrasport(null);
+ }
+ });
+ } catch (Exception e) {
+ LOG.error("Storm server experienced an exception while
creating a transport "
+ "using a JAAS principal context:"
+ e,
e);
diff --git
a/storm-client/src/jvm/org/apache/storm/security/auth/workertoken/WorkerTokenClientCallbackHandler.java
b/storm-client/src/jvm/org/apache/storm/security/auth/workertoken/WorkerTokenClientCallbackHandler.java
index 451c6d806..d85ff0382 100644
---
a/storm-client/src/jvm/org/apache/storm/security/auth/workertoken/WorkerTokenClientCallbackHandler.java
+++
b/storm-client/src/jvm/org/apache/storm/security/auth/workertoken/WorkerTokenClientCallbackHandler.java
@@ -12,12 +12,12 @@
package org.apache.storm.security.auth.workertoken;
-import java.security.AccessController;
import java.util.Base64;
import javax.security.auth.Subject;
import org.apache.storm.generated.WorkerToken;
import org.apache.storm.generated.WorkerTokenServiceType;
import org.apache.storm.security.auth.ClientAuthUtils;
+import org.apache.storm.security.auth.SubjectCompat;
import org.apache.storm.security.auth.ThriftConnectionType;
import org.apache.storm.security.auth.sasl.SimpleSaslClientCallbackHandler;
@@ -49,7 +49,7 @@ public class WorkerTokenClientCallbackHandler extends
SimpleSaslClientCallbackHa
WorkerTokenServiceType serviceType = type.getWtType();
WorkerToken ret = null;
if (serviceType != null) {
- Subject subject =
Subject.getSubject(AccessController.getContext());
+ Subject subject = SubjectCompat.currentSubject();
if (subject != null) {
ret = ClientAuthUtils.findWorkerToken(subject, serviceType);
}
diff --git a/storm-client/src/jvm/org/apache/storm/utils/HadoopLoginUtil.java
b/storm-client/src/jvm/org/apache/storm/utils/HadoopLoginUtil.java
index de8b417f6..76d3a1a67 100644
--- a/storm-client/src/jvm/org/apache/storm/utils/HadoopLoginUtil.java
+++ b/storm-client/src/jvm/org/apache/storm/utils/HadoopLoginUtil.java
@@ -20,11 +20,11 @@ package org.apache.storm.utils;
import java.lang.reflect.Method;
import java.net.UnknownHostException;
-import java.security.AccessController;
import java.security.PrivilegedAction;
import java.util.Map;
import javax.security.auth.Subject;
import org.apache.storm.Config;
+import org.apache.storm.security.auth.SubjectCompat;
import org.apache.storm.shade.com.google.common.base.Preconditions;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -136,7 +136,7 @@ public class HadoopLoginUtil {
Method doAsMethod = ugiClass.getMethod("doAs",
PrivilegedAction.class);
Object ugi = currentUserMethod.invoke(null);
return (Subject) doAsMethod.invoke(ugi,
- (PrivilegedAction<Subject>) () ->
Subject.getSubject(AccessController.getContext()));
+ (PrivilegedAction<Subject>) SubjectCompat::currentSubject);
} catch (Exception e) {
throw new RuntimeException("Error getting hadoop user!", e);
}
diff --git
a/storm-client/test/jvm/org/apache/storm/security/auth/SubjectCompatTest.java
b/storm-client/test/jvm/org/apache/storm/security/auth/SubjectCompatTest.java
new file mode 100644
index 000000000..f938eab4d
--- /dev/null
+++
b/storm-client/test/jvm/org/apache/storm/security/auth/SubjectCompatTest.java
@@ -0,0 +1,70 @@
+/**
+ * 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.storm.security.auth;
+
+import java.io.IOException;
+import javax.security.auth.Subject;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+
+class SubjectCompatTest {
+
+ @Test
+ void currentSubjectReturnsNullWhenNoneSet() {
+ assertNull(SubjectCompat.currentSubject());
+ }
+
+ @Test
+ void doAsExecutesCallableAndReturnsResult() throws Exception {
+ Subject subject = new Subject();
+ String result = SubjectCompat.doAs(subject, () -> "hello");
+ assertEquals("hello", result);
+ }
+
+ @Test
+ void doAsWithNullSubject() throws Exception {
+ String result = SubjectCompat.doAs(null, () -> "works");
+ assertEquals("works", result);
+ }
+
+ @Test
+ void doAsPropagatesCheckedException() {
+ Subject subject = new Subject();
+ assertThrows(IOException.class, () ->
+ SubjectCompat.doAs(subject, () -> {
+ throw new IOException("test error");
+ })
+ );
+ }
+
+ @Test
+ void doAsPropagatesRuntimeException() {
+ Subject subject = new Subject();
+ assertThrows(IllegalStateException.class, () ->
+ SubjectCompat.doAs(subject, () -> {
+ throw new IllegalStateException("test error");
+ })
+ );
+ }
+
+ @Test
+ void doAsSetsCurrentSubject() throws Exception {
+ Subject subject = new Subject();
+ Subject observed = SubjectCompat.doAs(subject,
SubjectCompat::currentSubject);
+ assertSame(subject, observed);
+ }
+}
diff --git
a/storm-server/src/test/java/org/apache/storm/security/auth/AuthTest.java
b/storm-server/src/test/java/org/apache/storm/security/auth/AuthTest.java
index 0706fe7c4..2720da049 100644
--- a/storm-server/src/test/java/org/apache/storm/security/auth/AuthTest.java
+++ b/storm-server/src/test/java/org/apache/storm/security/auth/AuthTest.java
@@ -16,8 +16,6 @@ import java.io.File;
import java.io.IOException;
import java.net.InetAddress;
import java.security.Principal;
-import java.security.PrivilegedActionException;
-import java.security.PrivilegedExceptionAction;
import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
@@ -184,8 +182,8 @@ public class AuthTest {
}
public static void tryConnectAs(Map<String, Object> conf, ThriftServer
server, Subject subject, String topoId)
- throws PrivilegedActionException {
- Subject.doAs(subject, (PrivilegedExceptionAction<Void>) () -> {
+ throws Exception {
+ SubjectCompat.doAs(subject, () -> {
try (NimbusClient client =
NimbusClient.Builder.withConf(conf).withTimeout(NIMBUS_TIMEOUT)
.buildWithNimbusHostPort("localhost", server.getPort())) {
client.getClient().activate(topoId); //Yes this should be a
topo name, but it makes this simpler...
@@ -195,7 +193,7 @@ public class AuthTest {
}
public static Subject testConnectWithTokenFor(WorkerTokenManager wtMan,
Map<String, Object> conf, ThriftServer server,
- String user, String topoId)
throws PrivilegedActionException {
+ String user, String topoId)
throws Exception {
WorkerToken wt =
wtMan.createOrUpdateTokenFor(WorkerTokenServiceType.NIMBUS, user, topoId);
Subject subject = createSubjectWith(wt);
tryConnectAs(conf, server, subject, topoId);