From c2dbdfb84a3a20a57814c0c7b7501db9a64477e2 Mon Sep 17 00:00:00 2001 From: Jongyoul Lee Date: Tue, 4 Aug 2026 01:08:59 +0900 Subject: [PATCH 1/2] Decouple zeppelin-server from Hadoop --- zeppelin-plugins/launcher/yarn/pom.xml | 8 +- .../zeppelin/interpreter/YarnAppMonitor.java | 100 +++++----- .../launcher/YarnProcessLaunchObserver.java | 58 ++++++ .../YarnRemoteInterpreterProcess.java | 1 + ...n.interpreter.remote.ProcessLaunchObserver | 16 ++ .../YarnProcessLaunchObserverTest.java | 56 ++++++ .../notebookrepo/filesystem/pom.xml | 14 +- .../zeppelin/healthcheck/HdfsHealthCheck.java | 2 +- .../recovery/FileSystemRecoveryStorage.java | 3 +- .../zeppelin/notebook/FileSystemStorage.java | 11 +- .../storage/FileSystemConfigStorage.java | 15 +- .../FileSystemRecoveryStorageTest.java | 67 +++++++ .../repo/FileSystemNotebookRepoTest.java | 6 +- .../storage/FileSystemConfigStorageTest.java | 67 +++++++ zeppelin-plugins/pom.xml | 6 +- zeppelin-plugins/security/hadoop/pom.xml | 84 ++++++++ ...adoopCredentialProviderSecretResolver.java | 43 ++++ .../realm/hadoop/HadoopGroupResolver.java | 42 ++++ .../KerberosAuthenticationFilter.java | 0 .../realm/kerberos/KerberosRealm.java | 31 ++- .../realm/kerberos/KerberosToken.java | 0 .../zeppelin/realm/kerberos/KerberosUtil.java | 0 ...pCredentialProviderSecretResolverTest.java | 50 +++++ zeppelin-server/pom.xml | 43 ++-- .../InterpreterSettingManager.java | 10 +- .../launcher/InterpreterLauncher.java | 12 +- .../launcher/StandardInterpreterLauncher.java | 21 +- .../interpreter/recovery/StopInterpreter.java | 11 +- .../remote/ExecRemoteInterpreterProcess.java | 94 ++++++--- .../remote/ProcessLaunchObserver.java | 34 ++++ .../notebook/repo/NotebookRepoSync.java | 17 +- .../apache/zeppelin/plugin/PluginManager.java | 188 ++++++++++++++++-- .../zeppelin/realm/ExternalLoginRealm.java | 49 +++++ .../apache/zeppelin/realm/GroupResolver.java | 26 +++ .../org/apache/zeppelin/realm/LdapRealm.java | 19 +- .../apache/zeppelin/realm/SecretResolver.java | 25 +++ .../realm/SecurityProviderLoader.java | 34 ++++ .../zeppelin/realm/ZeppelinRoleProvider.java | 25 +++ .../zeppelin/realm/jwt/KnoxJwtRealm.java | 114 +++++++---- .../apache/zeppelin/rest/LoginRestApi.java | 132 ++++-------- .../rest/message/ParagraphJobStatus.java | 2 +- .../zeppelin/server/ZeppelinServer.java | 139 ++++++++++++- .../service/ShiroAuthenticationService.java | 6 +- .../zeppelin/storage/ConfigStorage.java | 9 +- .../launcher/InterpreterLauncherTest.java | 30 +++ .../FileSystemRecoveryStorageTest.java | 118 ----------- .../zeppelin/realm/TestGroupResolver.java | 29 +++ .../zeppelin/realm/jwt/KnoxJwtRealmTest.java | 8 + .../zeppelin/recovery/RecoveryTest.java | 4 +- .../zeppelin/rest/AbstractTestRestApi.java | 1 + .../ZeppelinServerPluginClassLoadingTest.java | 126 ++++++++++++ 51 files changed, 1589 insertions(+), 417 deletions(-) rename {zeppelin-server => zeppelin-plugins/launcher/yarn}/src/main/java/org/apache/zeppelin/interpreter/YarnAppMonitor.java (52%) create mode 100644 zeppelin-plugins/launcher/yarn/src/main/java/org/apache/zeppelin/interpreter/launcher/YarnProcessLaunchObserver.java create mode 100644 zeppelin-plugins/launcher/yarn/src/main/resources/META-INF/services/org.apache.zeppelin.interpreter.remote.ProcessLaunchObserver create mode 100644 zeppelin-plugins/launcher/yarn/src/test/java/org/apache/zeppelin/interpreter/launcher/YarnProcessLaunchObserverTest.java rename {zeppelin-server => zeppelin-plugins/notebookrepo/filesystem}/src/main/java/org/apache/zeppelin/healthcheck/HdfsHealthCheck.java (98%) rename {zeppelin-server => zeppelin-plugins/notebookrepo/filesystem}/src/main/java/org/apache/zeppelin/interpreter/recovery/FileSystemRecoveryStorage.java (97%) rename {zeppelin-server => zeppelin-plugins/notebookrepo/filesystem}/src/main/java/org/apache/zeppelin/notebook/FileSystemStorage.java (95%) rename {zeppelin-server => zeppelin-plugins/notebookrepo/filesystem}/src/main/java/org/apache/zeppelin/storage/FileSystemConfigStorage.java (87%) create mode 100644 zeppelin-plugins/notebookrepo/filesystem/src/test/java/org/apache/zeppelin/interpreter/recovery/FileSystemRecoveryStorageTest.java create mode 100644 zeppelin-plugins/notebookrepo/filesystem/src/test/java/org/apache/zeppelin/storage/FileSystemConfigStorageTest.java create mode 100644 zeppelin-plugins/security/hadoop/pom.xml create mode 100644 zeppelin-plugins/security/hadoop/src/main/java/org/apache/zeppelin/realm/hadoop/HadoopCredentialProviderSecretResolver.java create mode 100644 zeppelin-plugins/security/hadoop/src/main/java/org/apache/zeppelin/realm/hadoop/HadoopGroupResolver.java rename {zeppelin-server => zeppelin-plugins/security/hadoop}/src/main/java/org/apache/zeppelin/realm/kerberos/KerberosAuthenticationFilter.java (100%) rename {zeppelin-server => zeppelin-plugins/security/hadoop}/src/main/java/org/apache/zeppelin/realm/kerberos/KerberosRealm.java (97%) rename {zeppelin-server => zeppelin-plugins/security/hadoop}/src/main/java/org/apache/zeppelin/realm/kerberos/KerberosToken.java (100%) rename {zeppelin-server => zeppelin-plugins/security/hadoop}/src/main/java/org/apache/zeppelin/realm/kerberos/KerberosUtil.java (100%) create mode 100644 zeppelin-plugins/security/hadoop/src/test/java/org/apache/zeppelin/realm/hadoop/HadoopCredentialProviderSecretResolverTest.java create mode 100644 zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/ProcessLaunchObserver.java create mode 100644 zeppelin-server/src/main/java/org/apache/zeppelin/realm/ExternalLoginRealm.java create mode 100644 zeppelin-server/src/main/java/org/apache/zeppelin/realm/GroupResolver.java create mode 100644 zeppelin-server/src/main/java/org/apache/zeppelin/realm/SecretResolver.java create mode 100644 zeppelin-server/src/main/java/org/apache/zeppelin/realm/SecurityProviderLoader.java create mode 100644 zeppelin-server/src/main/java/org/apache/zeppelin/realm/ZeppelinRoleProvider.java delete mode 100644 zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/recovery/FileSystemRecoveryStorageTest.java create mode 100644 zeppelin-server/src/test/java/org/apache/zeppelin/realm/TestGroupResolver.java create mode 100644 zeppelin-server/src/test/java/org/apache/zeppelin/server/ZeppelinServerPluginClassLoadingTest.java diff --git a/zeppelin-plugins/launcher/yarn/pom.xml b/zeppelin-plugins/launcher/yarn/pom.xml index e472227f2b5..9742139ad60 100644 --- a/zeppelin-plugins/launcher/yarn/pom.xml +++ b/zeppelin-plugins/launcher/yarn/pom.xml @@ -42,12 +42,18 @@ org.apache.hadoop hadoop-client-api - provided + compile org.apache.hadoop hadoop-client-runtime + compile + + + + org.slf4j + slf4j-api provided diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/YarnAppMonitor.java b/zeppelin-plugins/launcher/yarn/src/main/java/org/apache/zeppelin/interpreter/YarnAppMonitor.java similarity index 52% rename from zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/YarnAppMonitor.java rename to zeppelin-plugins/launcher/yarn/src/main/java/org/apache/zeppelin/interpreter/YarnAppMonitor.java index 48956b337be..8c346f6bb36 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/YarnAppMonitor.java +++ b/zeppelin-plugins/launcher/yarn/src/main/java/org/apache/zeppelin/interpreter/YarnAppMonitor.java @@ -17,6 +17,13 @@ package org.apache.zeppelin.interpreter; +import java.util.Iterator; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; + import org.apache.hadoop.yarn.api.records.ApplicationId; import org.apache.hadoop.yarn.api.records.ApplicationReport; import org.apache.hadoop.yarn.api.records.YarnApplicationState; @@ -29,16 +36,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import java.util.Iterator; -import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.Executors; -import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.TimeUnit; - -/** - * This class will launch a thread to check yarn app status regularly. - */ +/** This class launches a thread to check YARN application status regularly. */ public class YarnAppMonitor { private static final Logger LOGGER = LoggerFactory.getLogger(YarnAppMonitor.class); @@ -59,44 +57,23 @@ private YarnAppMonitor() { try { this.yarnClient = YarnClient.createYarnClient(); YarnConfiguration yarnConf = new YarnConfiguration(); - // disable timeline service as we only query yarn app here. - // Otherwise we may hit this kind of ERROR: - // java.lang.ClassNotFoundException: com.sun.jersey.api.client.config.ClientConfig + yarnConf.setClassLoader(YarnAppMonitor.class.getClassLoader()); + // Disable the timeline service because this client only queries application state. yarnConf.set("yarn.timeline-service.enabled", "false"); yarnClient.init(yarnConf); yarnClient.start(); - this.executor = Executors.newSingleThreadScheduledExecutor(new NamedThreadFactory("YarnAppsMonitor-Thread")); + this.executor = + Executors.newSingleThreadScheduledExecutor( + new NamedThreadFactory("YarnAppsMonitor-Thread")); this.apps = new ConcurrentHashMap<>(); - this.executor.scheduleAtFixedRate(() -> { - try { - Iterator> iter = apps.entrySet().iterator(); - while (iter.hasNext()) { - Map.Entry entry = iter.next(); - ApplicationId appId = entry.getKey(); - RemoteInterpreterManagedProcess interpreterManagedProcess = entry.getValue(); - ApplicationReport appReport = yarnClient.getApplicationReport(appId); - if (appReport.getYarnApplicationState() == YarnApplicationState.FAILED || - appReport.getYarnApplicationState() == YarnApplicationState.KILLED) { - String yarnDiagnostics = appReport.getDiagnostics(); - interpreterManagedProcess.processStopped("Yarn diagnostics: " + yarnDiagnostics); - iter.remove(); - LOGGER.info("Remove {} from YarnAppMonitor, because its state is {}", appId, - appReport.getYarnApplicationState()); - } else if (appReport.getYarnApplicationState() == YarnApplicationState.FINISHED) { - iter.remove(); - LOGGER.info("Remove {} from YarnAppMonitor, because its state is {}", appId, - appReport.getYarnApplicationState()); - } - } - } catch (Exception e) { - LOGGER.warn("Fail to check yarn app status", e); - } - }, - ZeppelinConfiguration - .getStaticInt(ConfVars.ZEPPELIN_INTERPRETER_YARN_MONITOR_INTERVAL_SECS), - ZeppelinConfiguration - .getStaticInt(ConfVars.ZEPPELIN_INTERPRETER_YARN_MONITOR_INTERVAL_SECS), - TimeUnit.SECONDS); + int monitorInterval = + ZeppelinConfiguration.getStaticInt( + ConfVars.ZEPPELIN_INTERPRETER_YARN_MONITOR_INTERVAL_SECS); + this.executor.scheduleAtFixedRate( + this::checkApplications, + monitorInterval, + monitorInterval, + TimeUnit.SECONDS); LOGGER.info("YarnAppMonitor is started"); } catch (Throwable e) { @@ -104,7 +81,40 @@ private YarnAppMonitor() { } } - public void addYarnApp(ApplicationId appId, RemoteInterpreterManagedProcess interpreterManagedProcess) { + private void checkApplications() { + try { + Iterator> iter = + apps.entrySet().iterator(); + while (iter.hasNext()) { + Map.Entry entry = iter.next(); + ApplicationId appId = entry.getKey(); + RemoteInterpreterManagedProcess interpreterManagedProcess = entry.getValue(); + ApplicationReport appReport = yarnClient.getApplicationReport(appId); + YarnApplicationState applicationState = appReport.getYarnApplicationState(); + if (applicationState == YarnApplicationState.FAILED + || applicationState == YarnApplicationState.KILLED) { + interpreterManagedProcess.processStopped( + "Yarn diagnostics: " + appReport.getDiagnostics()); + iter.remove(); + LOGGER.info( + "Remove {} from YarnAppMonitor, because its state is {}", + appId, + applicationState); + } else if (applicationState == YarnApplicationState.FINISHED) { + iter.remove(); + LOGGER.info( + "Remove {} from YarnAppMonitor, because its state is {}", + appId, + applicationState); + } + } + } catch (Exception e) { + LOGGER.warn("Fail to check yarn app status", e); + } + } + + public void addYarnApp( + ApplicationId appId, RemoteInterpreterManagedProcess interpreterManagedProcess) { LOGGER.info("Add {} to YarnAppMonitor", appId); this.apps.put(appId, interpreterManagedProcess); } diff --git a/zeppelin-plugins/launcher/yarn/src/main/java/org/apache/zeppelin/interpreter/launcher/YarnProcessLaunchObserver.java b/zeppelin-plugins/launcher/yarn/src/main/java/org/apache/zeppelin/interpreter/launcher/YarnProcessLaunchObserver.java new file mode 100644 index 00000000000..56e2389f0ca --- /dev/null +++ b/zeppelin-plugins/launcher/yarn/src/main/java/org/apache/zeppelin/interpreter/launcher/YarnProcessLaunchObserver.java @@ -0,0 +1,58 @@ +/* + * 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.zeppelin.interpreter.launcher; + +import java.util.function.BiConsumer; +import java.util.regex.Matcher; +import java.util.regex.Pattern; + +import org.apache.hadoop.yarn.api.records.ApplicationId; +import org.apache.zeppelin.interpreter.YarnAppMonitor; +import org.apache.zeppelin.interpreter.remote.ProcessLaunchObserver; +import org.apache.zeppelin.interpreter.remote.RemoteInterpreterManagedProcess; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** Detects an application submitted by a process launcher and registers it for YARN monitoring. */ +public class YarnProcessLaunchObserver implements ProcessLaunchObserver { + + private static final Logger LOGGER = LoggerFactory.getLogger(YarnProcessLaunchObserver.class); + private static final Pattern YARN_APP_PATTERN = Pattern.compile("Submitted application (\\w+)"); + + private final BiConsumer appConsumer; + + public YarnProcessLaunchObserver() { + this((appId, process) -> YarnAppMonitor.get().addYarnApp(appId, process)); + } + + YarnProcessLaunchObserver( + BiConsumer appConsumer) { + this.appConsumer = appConsumer; + } + + @Override + public void onProcessLaunch( + String launchOutput, RemoteInterpreterManagedProcess interpreterProcess) { + Matcher matcher = YARN_APP_PATTERN.matcher(launchOutput); + if (matcher.find()) { + String appId = matcher.group(1); + LOGGER.info("Detected yarn app: {}, add it to YarnAppMonitor", appId); + appConsumer.accept(ApplicationId.fromString(appId), interpreterProcess); + } + } +} diff --git a/zeppelin-plugins/launcher/yarn/src/main/java/org/apache/zeppelin/interpreter/launcher/YarnRemoteInterpreterProcess.java b/zeppelin-plugins/launcher/yarn/src/main/java/org/apache/zeppelin/interpreter/launcher/YarnRemoteInterpreterProcess.java index f47eaf68d76..4281d6bd173 100644 --- a/zeppelin-plugins/launcher/yarn/src/main/java/org/apache/zeppelin/interpreter/launcher/YarnRemoteInterpreterProcess.java +++ b/zeppelin-plugins/launcher/yarn/src/main/java/org/apache/zeppelin/interpreter/launcher/YarnRemoteInterpreterProcess.java @@ -113,6 +113,7 @@ public YarnRemoteInterpreterProcess( this.properties = properties; this.envs = envs; this.hadoopConf = new YarnConfiguration(); + this.hadoopConf.setClassLoader(YarnRemoteInterpreterProcess.class.getClassLoader()); // Add core-site.xml and yarn-site.xml. This is for integration test where using MiniHadoopCluster. if (properties.containsKey("HADOOP_CONF_DIR") && !StringUtils.isBlank(properties.getProperty("HADOOP_CONF_DIR"))) { diff --git a/zeppelin-plugins/launcher/yarn/src/main/resources/META-INF/services/org.apache.zeppelin.interpreter.remote.ProcessLaunchObserver b/zeppelin-plugins/launcher/yarn/src/main/resources/META-INF/services/org.apache.zeppelin.interpreter.remote.ProcessLaunchObserver new file mode 100644 index 00000000000..de042000e1e --- /dev/null +++ b/zeppelin-plugins/launcher/yarn/src/main/resources/META-INF/services/org.apache.zeppelin.interpreter.remote.ProcessLaunchObserver @@ -0,0 +1,16 @@ +# 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. + +org.apache.zeppelin.interpreter.launcher.YarnProcessLaunchObserver diff --git a/zeppelin-plugins/launcher/yarn/src/test/java/org/apache/zeppelin/interpreter/launcher/YarnProcessLaunchObserverTest.java b/zeppelin-plugins/launcher/yarn/src/test/java/org/apache/zeppelin/interpreter/launcher/YarnProcessLaunchObserverTest.java new file mode 100644 index 00000000000..c877cf0b5f2 --- /dev/null +++ b/zeppelin-plugins/launcher/yarn/src/test/java/org/apache/zeppelin/interpreter/launcher/YarnProcessLaunchObserverTest.java @@ -0,0 +1,56 @@ +/* + * 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.zeppelin.interpreter.launcher; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.mockito.Mockito.mock; + +import java.util.concurrent.atomic.AtomicReference; + +import org.apache.hadoop.yarn.api.records.ApplicationId; +import org.apache.zeppelin.interpreter.remote.RemoteInterpreterManagedProcess; +import org.junit.jupiter.api.Test; + +class YarnProcessLaunchObserverTest { + + @Test + void detectsSubmittedYarnApplication() { + AtomicReference detectedApp = new AtomicReference<>(); + RemoteInterpreterManagedProcess process = mock(RemoteInterpreterManagedProcess.class); + YarnProcessLaunchObserver observer = + new YarnProcessLaunchObserver((appId, ignored) -> detectedApp.set(appId)); + + observer.onProcessLaunch( + "INFO Client: Submitted application application_1720000000000_0042", process); + + assertEquals("application_1720000000000_0042", detectedApp.get().toString()); + } + + @Test + void ignoresLaunchOutputWithoutSubmittedApplication() { + AtomicReference detectedApp = new AtomicReference<>(); + YarnProcessLaunchObserver observer = + new YarnProcessLaunchObserver((appId, ignored) -> detectedApp.set(appId)); + + observer.onProcessLaunch( + "INFO Client: Application report for application_1720000000000_0042", null); + + assertNull(detectedApp.get()); + } +} diff --git a/zeppelin-plugins/notebookrepo/filesystem/pom.xml b/zeppelin-plugins/notebookrepo/filesystem/pom.xml index b369a3959aa..2d838840920 100644 --- a/zeppelin-plugins/notebookrepo/filesystem/pom.xml +++ b/zeppelin-plugins/notebookrepo/filesystem/pom.xml @@ -30,8 +30,8 @@ notebookrepo-filesystem jar - Zeppelin: Plugin FileSystemNotebookRepo - NotebookRepo implementation based on Hadoop FileSystem + Zeppelin: Plugin Hadoop FileSystem Storage + Notebook, configuration, and recovery storage based on Hadoop FileSystem 2.1.4 @@ -39,9 +39,19 @@ + + org.apache.hadoop + hadoop-client-api + compile + org.apache.hadoop hadoop-client-runtime + runtime + + + org.slf4j + slf4j-api provided diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/healthcheck/HdfsHealthCheck.java b/zeppelin-plugins/notebookrepo/filesystem/src/main/java/org/apache/zeppelin/healthcheck/HdfsHealthCheck.java similarity index 98% rename from zeppelin-server/src/main/java/org/apache/zeppelin/healthcheck/HdfsHealthCheck.java rename to zeppelin-plugins/notebookrepo/filesystem/src/main/java/org/apache/zeppelin/healthcheck/HdfsHealthCheck.java index 5f7f48afd5c..ad9a95a5edf 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/healthcheck/HdfsHealthCheck.java +++ b/zeppelin-plugins/notebookrepo/filesystem/src/main/java/org/apache/zeppelin/healthcheck/HdfsHealthCheck.java @@ -34,7 +34,7 @@ public class HdfsHealthCheck extends HealthCheck { */ public HdfsHealthCheck(FileSystemStorage fs, Path path) { this.fs = fs; - this.path= path; + this.path = path; } @Override protected Result check() throws Exception { diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/recovery/FileSystemRecoveryStorage.java b/zeppelin-plugins/notebookrepo/filesystem/src/main/java/org/apache/zeppelin/interpreter/recovery/FileSystemRecoveryStorage.java similarity index 97% rename from zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/recovery/FileSystemRecoveryStorage.java rename to zeppelin-plugins/notebookrepo/filesystem/src/main/java/org/apache/zeppelin/interpreter/recovery/FileSystemRecoveryStorage.java index 26cecb256ed..59afae1eabf 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/recovery/FileSystemRecoveryStorage.java +++ b/zeppelin-plugins/notebookrepo/filesystem/src/main/java/org/apache/zeppelin/interpreter/recovery/FileSystemRecoveryStorage.java @@ -53,7 +53,8 @@ public FileSystemRecoveryStorage(ZeppelinConfiguration zConf, throws IOException { super(zConf); this.interpreterSettingManager = interpreterSettingManager; - String recoveryDirProperty = zConf.getString(ZeppelinConfiguration.ConfVars.ZEPPELIN_RECOVERY_DIR); + String recoveryDirProperty = + zConf.getString(ZeppelinConfiguration.ConfVars.ZEPPELIN_RECOVERY_DIR); this.fs = new FileSystemStorage(zConf, recoveryDirProperty); LOGGER.info("Creating FileSystem: " + this.fs.getFs().getClass().getName() + " for Zeppelin Recovery."); diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/notebook/FileSystemStorage.java b/zeppelin-plugins/notebookrepo/filesystem/src/main/java/org/apache/zeppelin/notebook/FileSystemStorage.java similarity index 95% rename from zeppelin-server/src/main/java/org/apache/zeppelin/notebook/FileSystemStorage.java rename to zeppelin-plugins/notebookrepo/filesystem/src/main/java/org/apache/zeppelin/notebook/FileSystemStorage.java index 5fc60e74e55..f78f3d1198c 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/notebook/FileSystemStorage.java +++ b/zeppelin-plugins/notebookrepo/filesystem/src/main/java/org/apache/zeppelin/notebook/FileSystemStorage.java @@ -82,6 +82,7 @@ public class FileSystemStorage { public FileSystemStorage(ZeppelinConfiguration zConf, String path) throws IOException { this.zConf = zConf; this.hadoopConf = new Configuration(); + this.hadoopConf.setClassLoader(FileSystemStorage.class.getClassLoader()); URI zepConfigURI; URI defaultFSURI; @@ -169,7 +170,7 @@ public List call() throws IOException { }); } - // recursive search path, (TODO zjffdu, list folder in sub folder on demand, instead of load all + // recursive search path, (TODO(zjffdu): list folder in sub folder on demand, instead of load all // data when zeppelin server start) public List listAll(final Path path) throws IOException { return callHdfsOperation(new HdfsOperation>() { @@ -213,17 +214,19 @@ public String call() throws IOException { LOGGER.debug("Read from file: {}", file); ByteArrayOutputStream noteBytes = new ByteArrayOutputStream(); IOUtils.copyBytes(fs.open(file), noteBytes, hadoopConf); - return noteBytes.toString(zConf.getString(ZeppelinConfiguration.ConfVars.ZEPPELIN_ENCODING)); + return noteBytes.toString( + zConf.getString(ZeppelinConfiguration.ConfVars.ZEPPELIN_ENCODING)); } }); } public void writeFile(final String content, final Path file, boolean writeTempFileFirst) throws IOException { - writeFile(content, file, writeTempFileFirst, null); + writeFile(content, file, writeTempFileFirst, null); } - public void writeFile(final String content, final Path file, boolean writeTempFileFirst, Set permissions) + public void writeFile(final String content, final Path file, boolean writeTempFileFirst, + Set permissions) throws IOException { FsPermission fsPermission; if (permissions == null || permissions.isEmpty()) { diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/storage/FileSystemConfigStorage.java b/zeppelin-plugins/notebookrepo/filesystem/src/main/java/org/apache/zeppelin/storage/FileSystemConfigStorage.java similarity index 87% rename from zeppelin-server/src/main/java/org/apache/zeppelin/storage/FileSystemConfigStorage.java rename to zeppelin-plugins/notebookrepo/filesystem/src/main/java/org/apache/zeppelin/storage/FileSystemConfigStorage.java index ac7d108f27b..cf7c5ab6364 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/storage/FileSystemConfigStorage.java +++ b/zeppelin-plugins/notebookrepo/filesystem/src/main/java/org/apache/zeppelin/storage/FileSystemConfigStorage.java @@ -51,14 +51,18 @@ public FileSystemConfigStorage(ZeppelinConfiguration zConf) throws IOException { super(zConf); String configDir = zConf.getConfigFSDir(false); this.fs = new FileSystemStorage(zConf, configDir); - LOGGER.info("Creating FileSystem: {} for Zeppelin Config", this.fs.getFs().getClass().getName()); + LOGGER.info("Creating FileSystem: {} for Zeppelin Config", + this.fs.getFs().getClass().getName()); Path configPath = this.fs.makeQualified(new Path(configDir)); this.fs.tryMkDir(configPath); LOGGER.info("Using folder {} to store Zeppelin Config", configPath); - this.interpreterSettingPath = fs.makeQualified(new Path(zConf.getInterpreterSettingPath(false))); - this.authorizationPath = fs.makeQualified(new Path(zConf.getNotebookAuthorizationPath(false))); + this.interpreterSettingPath = + fs.makeQualified(new Path(zConf.getInterpreterSettingPath(false))); + this.authorizationPath = + fs.makeQualified(new Path(zConf.getNotebookAuthorizationPath(false))); this.credentialPath = fs.makeQualified(new Path(zConf.getCredentialsPath(false))); - HealthChecks.getHealthCheckLivenessRegistry().register(STORAGE_HEALTHCHECK_NAME, new HdfsHealthCheck(this.fs, configPath)); + HealthChecks.getHealthCheckLivenessRegistry().register( + STORAGE_HEALTHCHECK_NAME, new HdfsHealthCheck(this.fs, configPath)); } @Override @@ -108,7 +112,8 @@ public String loadCredentials() throws IOException { @Override public void saveCredentials(String credentials) throws IOException { LOGGER.info("Save Credentials to file: {}", credentialPath); - Set permissions = EnumSet.of(PosixFilePermission.OWNER_READ, PosixFilePermission.OWNER_WRITE); + Set permissions = + EnumSet.of(PosixFilePermission.OWNER_READ, PosixFilePermission.OWNER_WRITE); fs.writeFile(credentials, credentialPath, false, permissions); } diff --git a/zeppelin-plugins/notebookrepo/filesystem/src/test/java/org/apache/zeppelin/interpreter/recovery/FileSystemRecoveryStorageTest.java b/zeppelin-plugins/notebookrepo/filesystem/src/test/java/org/apache/zeppelin/interpreter/recovery/FileSystemRecoveryStorageTest.java new file mode 100644 index 00000000000..4a06f3c246a --- /dev/null +++ b/zeppelin-plugins/notebookrepo/filesystem/src/test/java/org/apache/zeppelin/interpreter/recovery/FileSystemRecoveryStorageTest.java @@ -0,0 +1,67 @@ +/* + * 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.zeppelin.interpreter.recovery; + +import org.apache.zeppelin.conf.ZeppelinConfiguration; +import org.apache.zeppelin.interpreter.InterpreterSetting; +import org.apache.zeppelin.interpreter.InterpreterSettingManager; +import org.apache.zeppelin.interpreter.launcher.InterpreterClient; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.Collections; +import java.util.Properties; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +class FileSystemRecoveryStorageTest { + + @TempDir + Path recoveryDir; + + @Test + void persistsRecoveryDataThroughHadoopFileSystem() throws Exception { + ZeppelinConfiguration zConf = ZeppelinConfiguration.load(); + zConf.setProperty(ZeppelinConfiguration.ConfVars.ZEPPELIN_RECOVERY_DIR.getVarName(), + recoveryDir.toString()); + + InterpreterSetting setting = mock(InterpreterSetting.class); + when(setting.getAllInterpreterGroups()).thenReturn(Collections.emptyList()); + when(setting.getJavaProperties()).thenReturn(new Properties()); + + InterpreterSettingManager settingManager = mock(InterpreterSettingManager.class); + when(settingManager.getInterpreterSettingByName("test")).thenReturn(setting); + when(settingManager.getByName("test")).thenReturn(setting); + + InterpreterClient client = mock(InterpreterClient.class); + when(client.getInterpreterSettingName()).thenReturn("test"); + + FileSystemRecoveryStorage storage = new FileSystemRecoveryStorage(zConf, settingManager); + storage.onInterpreterClientStart(client); + + Path recoveryFile = recoveryDir.resolve("test.recovery"); + assertTrue(Files.exists(recoveryFile)); + assertEquals("", Files.readString(recoveryFile)); + assertTrue(storage.restore().isEmpty()); + } +} diff --git a/zeppelin-plugins/notebookrepo/filesystem/src/test/java/org/apache/zeppelin/notebook/repo/FileSystemNotebookRepoTest.java b/zeppelin-plugins/notebookrepo/filesystem/src/test/java/org/apache/zeppelin/notebook/repo/FileSystemNotebookRepoTest.java index 1d7f8d19def..91fccb38f44 100644 --- a/zeppelin-plugins/notebookrepo/filesystem/src/test/java/org/apache/zeppelin/notebook/repo/FileSystemNotebookRepoTest.java +++ b/zeppelin-plugins/notebookrepo/filesystem/src/test/java/org/apache/zeppelin/notebook/repo/FileSystemNotebookRepoTest.java @@ -53,7 +53,8 @@ class FileSystemNotebookRepoTest { @BeforeEach void setUp() throws IOException { - notebookDir = Files.createTempDirectory("FileSystemNotebookRepoTest").toFile().getAbsolutePath(); + notebookDir = Files.createTempDirectory("FileSystemNotebookRepoTest") + .toFile().getAbsolutePath(); zConf = ZeppelinConfiguration.load(); noteParser = new GsonNoteParser(zConf); zConf.setProperty(ZeppelinConfiguration.ConfVars.ZEPPELIN_NOTEBOOK_DIR.getVarName(), @@ -131,7 +132,8 @@ void testBasics() throws IOException { @Test void testComplicatedScenarios() throws IOException { - // scenario_1: notebook_dir is not clean. There're some unrecognized dir and file under notebook_dir + // scenario_1: notebook_dir is not clean. There're some unrecognized dir and file under + // notebook_dir fs.mkdirs(new Path(notebookDir, "1/2")); OutputStream out = fs.create(new Path(notebookDir, "1/a.json")); out.close(); diff --git a/zeppelin-plugins/notebookrepo/filesystem/src/test/java/org/apache/zeppelin/storage/FileSystemConfigStorageTest.java b/zeppelin-plugins/notebookrepo/filesystem/src/test/java/org/apache/zeppelin/storage/FileSystemConfigStorageTest.java new file mode 100644 index 00000000000..14c1e7aeace --- /dev/null +++ b/zeppelin-plugins/notebookrepo/filesystem/src/test/java/org/apache/zeppelin/storage/FileSystemConfigStorageTest.java @@ -0,0 +1,67 @@ +/* + * 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.zeppelin.storage; + +import org.apache.zeppelin.conf.ZeppelinConfiguration; +import org.apache.zeppelin.healthcheck.HealthChecks; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +import java.net.URL; +import java.net.URLClassLoader; +import java.nio.file.Path; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class FileSystemConfigStorageTest { + + @TempDir + Path configDir; + + @AfterEach + void removeHealthCheck() { + HealthChecks.getHealthCheckLivenessRegistry().unregister( + ConfigStorage.STORAGE_HEALTHCHECK_NAME); + } + + @Test + void persistsCredentialsAndRegistersAHealthyFileSystem() throws Exception { + ZeppelinConfiguration zConf = ZeppelinConfiguration.load(); + zConf.setProperty(ZeppelinConfiguration.ConfVars.ZEPPELIN_CONFIG_FS_DIR.getVarName(), + configDir.toString()); + zConf.setProperty(ZeppelinConfiguration.ConfVars.ZEPPELIN_CONFIG_STORAGE_CLASS.getVarName(), + FileSystemConfigStorage.class.getName()); + + FileSystemConfigStorage storage; + Thread thread = Thread.currentThread(); + ClassLoader previousClassLoader = thread.getContextClassLoader(); + try (URLClassLoader emptyClassLoader = new URLClassLoader(new URL[0], null)) { + thread.setContextClassLoader(emptyClassLoader); + storage = new FileSystemConfigStorage(zConf); + storage.saveCredentials("secret"); + } finally { + thread.setContextClassLoader(previousClassLoader); + } + + assertEquals("secret", storage.loadCredentials()); + assertTrue(HealthChecks.getHealthCheckLivenessRegistry() + .runHealthCheck(ConfigStorage.STORAGE_HEALTHCHECK_NAME).isHealthy()); + } +} diff --git a/zeppelin-plugins/pom.xml b/zeppelin-plugins/pom.xml index ac859b3c764..8729c4bcac7 100644 --- a/zeppelin-plugins/pom.xml +++ b/zeppelin-plugins/pom.xml @@ -31,11 +31,6 @@ pom Zeppelin: Plugins Parent Zeppelin Plugins Parent - - - provided - - notebookrepo/s3 notebookrepo/github @@ -49,6 +44,7 @@ launcher/docker launcher/yarn launcher/flink + security/hadoop diff --git a/zeppelin-plugins/security/hadoop/pom.xml b/zeppelin-plugins/security/hadoop/pom.xml new file mode 100644 index 00000000000..556e130bdef --- /dev/null +++ b/zeppelin-plugins/security/hadoop/pom.xml @@ -0,0 +1,84 @@ + + + + + + 4.0.0 + + + zengine-plugins-parent + org.apache.zeppelin + 0.13.0-SNAPSHOT + ../.. + + + security-hadoop + jar + Zeppelin: Hadoop Security Plugin + Optional Hadoop-backed Kerberos, group and credential provider integrations + + + Security/Hadoop + 2.0.0-M15 + + + + + org.apache.hadoop + hadoop-client-api + compile + + + + org.apache.hadoop + hadoop-client-runtime + compile + + + + org.apache.directory.server + apacheds-kerberos-codec + ${kerberos.version} + + + + org.slf4j + slf4j-api + provided + + + + + + + maven-dependency-plugin + + + copy-plugin-dependencies + + + slf4j-api + + + + + + + diff --git a/zeppelin-plugins/security/hadoop/src/main/java/org/apache/zeppelin/realm/hadoop/HadoopCredentialProviderSecretResolver.java b/zeppelin-plugins/security/hadoop/src/main/java/org/apache/zeppelin/realm/hadoop/HadoopCredentialProviderSecretResolver.java new file mode 100644 index 00000000000..cc9f97cace5 --- /dev/null +++ b/zeppelin-plugins/security/hadoop/src/main/java/org/apache/zeppelin/realm/hadoop/HadoopCredentialProviderSecretResolver.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.zeppelin.realm.hadoop; + +import java.io.IOException; +import java.util.List; + +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.security.alias.CredentialProvider; +import org.apache.hadoop.security.alias.CredentialProviderFactory; +import org.apache.zeppelin.realm.SecretResolver; + +/** Reads secrets from a Hadoop credential provider such as a JCEKS file. */ +public class HadoopCredentialProviderSecretResolver implements SecretResolver { + + @Override + public char[] resolve(String providerPath, String alias) throws IOException { + Configuration configuration = new Configuration(); + configuration.setClassLoader( + HadoopCredentialProviderSecretResolver.class.getClassLoader()); + configuration.set(CredentialProviderFactory.CREDENTIAL_PROVIDER_PATH, providerPath); + List providers = CredentialProviderFactory.getProviders(configuration); + if (providers.isEmpty()) { + return null; + } + CredentialProvider.CredentialEntry entry = providers.get(0).getCredentialEntry(alias); + return entry == null ? null : entry.getCredential(); + } +} diff --git a/zeppelin-plugins/security/hadoop/src/main/java/org/apache/zeppelin/realm/hadoop/HadoopGroupResolver.java b/zeppelin-plugins/security/hadoop/src/main/java/org/apache/zeppelin/realm/hadoop/HadoopGroupResolver.java new file mode 100644 index 00000000000..0177ab35e12 --- /dev/null +++ b/zeppelin-plugins/security/hadoop/src/main/java/org/apache/zeppelin/realm/hadoop/HadoopGroupResolver.java @@ -0,0 +1,42 @@ +/* + * 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.zeppelin.realm.hadoop; + +import java.io.IOException; +import java.util.HashSet; +import java.util.Set; + +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.security.Groups; +import org.apache.zeppelin.realm.GroupResolver; + +/** Hadoop-backed group resolver for realms that support Hadoop group mappings. */ +public class HadoopGroupResolver implements GroupResolver { + + private final Groups groups; + + public HadoopGroupResolver() { + Configuration configuration = new Configuration(); + configuration.setClassLoader(HadoopGroupResolver.class.getClassLoader()); + groups = new Groups(configuration); + } + + @Override + public Set resolve(String principal) throws IOException { + return new HashSet<>(groups.getGroups(principal)); + } +} diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/realm/kerberos/KerberosAuthenticationFilter.java b/zeppelin-plugins/security/hadoop/src/main/java/org/apache/zeppelin/realm/kerberos/KerberosAuthenticationFilter.java similarity index 100% rename from zeppelin-server/src/main/java/org/apache/zeppelin/realm/kerberos/KerberosAuthenticationFilter.java rename to zeppelin-plugins/security/hadoop/src/main/java/org/apache/zeppelin/realm/kerberos/KerberosAuthenticationFilter.java diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/realm/kerberos/KerberosRealm.java b/zeppelin-plugins/security/hadoop/src/main/java/org/apache/zeppelin/realm/kerberos/KerberosRealm.java similarity index 97% rename from zeppelin-server/src/main/java/org/apache/zeppelin/realm/kerberos/KerberosRealm.java rename to zeppelin-plugins/security/hadoop/src/main/java/org/apache/zeppelin/realm/kerberos/KerberosRealm.java index 0a09853aa7f..3bf2fa49a57 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/realm/kerberos/KerberosRealm.java +++ b/zeppelin-plugins/security/hadoop/src/main/java/org/apache/zeppelin/realm/kerberos/KerberosRealm.java @@ -34,6 +34,7 @@ import org.apache.shiro.authz.SimpleAuthorizationInfo; import org.apache.shiro.realm.AuthorizingRealm; import org.apache.shiro.subject.PrincipalCollection; +import org.apache.zeppelin.realm.ExternalLoginRealm; import org.ietf.jgss.GSSException; import org.ietf.jgss.GSSContext; import org.ietf.jgss.GSSCredential; @@ -85,7 +86,7 @@ * authc = org.apache.zeppelin.realm.kerberos.KerberosAuthenticationFilter * */ -public class KerberosRealm extends AuthorizingRealm { +public class KerberosRealm extends AuthorizingRealm implements ExternalLoginRealm { private static final Logger LOGGER = LoggerFactory.getLogger(KerberosRealm.class); // Configs to set in shiro.ini @@ -203,6 +204,28 @@ public boolean supports(org.apache.shiro.authc.AuthenticationToken token) { return token instanceof KerberosToken; } + @Override + public org.apache.shiro.authc.AuthenticationToken getLoginAuthenticationToken( + Map cookies) + throws org.apache.shiro.authc.AuthenticationException { + return getKerberosTokenFromCookies(cookies); + } + + @Override + public String getLoginPrincipal(org.apache.shiro.authc.AuthenticationToken token) { + return (String) token.getPrincipal(); + } + + @Override + public boolean shouldRedirectOnMissingToken() { + return false; + } + + @Override + public int getLoginPriority() { + return 0; + } + /** * Initializes the KerberosRealm by 'kinit'ing using principal and keytab. *

@@ -276,6 +299,7 @@ protected void onInit() { } Configuration hadoopConfig = new Configuration(); + hadoopConfig.setClassLoader(KerberosRealm.class.getClassLoader()); hadoopGroups = new Groups(hadoopConfig); } catch (Exception ex) { @@ -984,6 +1008,11 @@ public String getLogout() { return logout; } + @Override + public String getLogin() { + return null; + } + public void setLogout(String logout) { this.logout = logout; } diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/realm/kerberos/KerberosToken.java b/zeppelin-plugins/security/hadoop/src/main/java/org/apache/zeppelin/realm/kerberos/KerberosToken.java similarity index 100% rename from zeppelin-server/src/main/java/org/apache/zeppelin/realm/kerberos/KerberosToken.java rename to zeppelin-plugins/security/hadoop/src/main/java/org/apache/zeppelin/realm/kerberos/KerberosToken.java diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/realm/kerberos/KerberosUtil.java b/zeppelin-plugins/security/hadoop/src/main/java/org/apache/zeppelin/realm/kerberos/KerberosUtil.java similarity index 100% rename from zeppelin-server/src/main/java/org/apache/zeppelin/realm/kerberos/KerberosUtil.java rename to zeppelin-plugins/security/hadoop/src/main/java/org/apache/zeppelin/realm/kerberos/KerberosUtil.java diff --git a/zeppelin-plugins/security/hadoop/src/test/java/org/apache/zeppelin/realm/hadoop/HadoopCredentialProviderSecretResolverTest.java b/zeppelin-plugins/security/hadoop/src/test/java/org/apache/zeppelin/realm/hadoop/HadoopCredentialProviderSecretResolverTest.java new file mode 100644 index 00000000000..2f263f178f6 --- /dev/null +++ b/zeppelin-plugins/security/hadoop/src/test/java/org/apache/zeppelin/realm/hadoop/HadoopCredentialProviderSecretResolverTest.java @@ -0,0 +1,50 @@ +/* + * 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.zeppelin.realm.hadoop; + +import static org.junit.jupiter.api.Assertions.assertArrayEquals; + +import java.nio.file.Path; + +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.security.alias.CredentialProvider; +import org.apache.hadoop.security.alias.CredentialProviderFactory; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +class HadoopCredentialProviderSecretResolverTest { + + @TempDir + Path tempDir; + + @Test + void resolvesExistingJceksAlias() throws Exception { + String providerPath = "jceks://file" + tempDir.resolve("zeppelin.jceks").toAbsolutePath(); + Configuration configuration = new Configuration(); + configuration.set(CredentialProviderFactory.CREDENTIAL_PROVIDER_PATH, providerPath); + CredentialProvider provider = CredentialProviderFactory.getProviders(configuration).get(0); + provider.createCredentialEntry("ldapRealm.systemPassword", "secret".toCharArray()); + provider.flush(); + + HadoopCredentialProviderSecretResolver resolver = + new HadoopCredentialProviderSecretResolver(); + + assertArrayEquals( + "secret".toCharArray(), + resolver.resolve(providerPath, "ldapRealm.systemPassword")); + } +} diff --git a/zeppelin-server/pom.xml b/zeppelin-server/pom.xml index 8dbdb575673..aa2ac79b281 100644 --- a/zeppelin-server/pom.xml +++ b/zeppelin-server/pom.xml @@ -39,7 +39,6 @@ 1.11 4.1.0 9.37.4 - 2.0.0-M15 32.0.0-jre 8.7.0 1.18.0 @@ -69,6 +68,14 @@ zeppelin-interpreter ${project.version} + + org.apache.hadoop + hadoop-client-api + + + org.apache.hadoop + hadoop-client-runtime + javax.inject javax.inject @@ -219,7 +226,7 @@ commons-vfs2 ${commons.vfs2.version} - + org.apache.hadoop hadoop-hdfs-client @@ -444,12 +451,6 @@ ${quartz.scheduler.version} - - org.apache.directory.server - apacheds-kerberos-codec - ${kerberos.version} - - org.apache.zeppelin zeppelin-test @@ -551,11 +552,6 @@ 1.1 - - org.apache.hadoop - hadoop-client-runtime - ${hadoop.deps.scope} - @@ -644,6 +640,27 @@ true + + enforce-no-hadoop-dependencies + + enforce + + + + + + org.apache.hadoop:* + + true + + Hadoop integrations belong in optional plugins and must not appear on the + zeppelin-server classpath. + + + + true + + diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSettingManager.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSettingManager.java index e2959382206..c8683c26917 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSettingManager.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSettingManager.java @@ -63,7 +63,6 @@ import org.apache.zeppelin.resource.ResourceSet; import org.apache.zeppelin.scheduler.Job; import org.apache.zeppelin.user.AuthenticationInfo; -import org.apache.zeppelin.util.ReflectionUtils; import org.apache.zeppelin.storage.ConfigStorage; import org.eclipse.jetty.util.annotation.ManagedAttribute; import org.eclipse.jetty.util.annotation.ManagedObject; @@ -196,11 +195,10 @@ public InterpreterSettingManager(ZeppelinConfiguration zConf, this.interpreterEventServer = new RemoteInterpreterEventServer(zConf, this); this.interpreterEventServer.start(); - this.recoveryStorage = - ReflectionUtils.createClazzInstance( - zConf.getRecoveryStorageClass(), - new Class[] {ZeppelinConfiguration.class, InterpreterSettingManager.class}, - new Object[] {zConf, this}); + this.recoveryStorage = pluginManager.createPluginInstance( + zConf.getRecoveryStorageClass(), + new Class[] {ZeppelinConfiguration.class, InterpreterSettingManager.class}, + new Object[] {zConf, this}); LOGGER.info("Using RecoveryStorage: {}", this.recoveryStorage.getClass().getName()); diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/launcher/InterpreterLauncher.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/launcher/InterpreterLauncher.java index 7190bea4896..dbbef8a3597 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/launcher/InterpreterLauncher.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/launcher/InterpreterLauncher.java @@ -102,8 +102,16 @@ public InterpreterClient launch(InterpreterLaunchContext context) throws IOExcep } } - // launch it via sub class implementation without recovering. - return launchDirectly(context); + // Launch external plugins with their own classloader as TCCL. Libraries such as Hadoop use + // TCCL for configuration resources and service discovery. + Thread thread = Thread.currentThread(); + ClassLoader previousClassLoader = thread.getContextClassLoader(); + try { + thread.setContextClassLoader(getClass().getClassLoader()); + return launchDirectly(context); + } finally { + thread.setContextClassLoader(previousClassLoader); + } } /** diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/launcher/StandardInterpreterLauncher.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/launcher/StandardInterpreterLauncher.java index 77e6e7bddce..fbacbda37b6 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/launcher/StandardInterpreterLauncher.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/launcher/StandardInterpreterLauncher.java @@ -25,6 +25,7 @@ import org.apache.zeppelin.interpreter.InterpreterRunner; import org.apache.zeppelin.interpreter.recovery.RecoveryStorage; import org.apache.zeppelin.interpreter.remote.ExecRemoteInterpreterProcess; +import org.apache.zeppelin.interpreter.remote.ProcessLaunchObserver; import org.apache.zeppelin.interpreter.remote.RemoteInterpreterRunningProcess; import org.apache.zeppelin.interpreter.remote.RemoteInterpreterUtils; import org.slf4j.Logger; @@ -32,6 +33,7 @@ import java.io.File; import java.io.IOException; +import java.util.List; import java.util.Map; /** @@ -40,11 +42,27 @@ public class StandardInterpreterLauncher extends InterpreterLauncher { private static final Logger LOGGER = LoggerFactory.getLogger(StandardInterpreterLauncher.class); + private ProcessLaunchObserver processLaunchObserver = ProcessLaunchObserver.NO_OP; public StandardInterpreterLauncher(ZeppelinConfiguration zConf, RecoveryStorage recoveryStorage) { super(zConf, recoveryStorage); } + public void setProcessLaunchObservers(List observers) { + this.processLaunchObserver = (launchOutput, interpreterProcess) -> { + for (ProcessLaunchObserver observer : observers) { + Thread thread = Thread.currentThread(); + ClassLoader previousClassLoader = thread.getContextClassLoader(); + try { + thread.setContextClassLoader(observer.getClass().getClassLoader()); + observer.onProcessLaunch(launchOutput, interpreterProcess); + } finally { + thread.setContextClassLoader(previousClassLoader); + } + } + }; + } + @Override public InterpreterClient launchDirectly(InterpreterLaunchContext context) throws IOException { LOGGER.info("Launching new interpreter process of {}", context.getInterpreterSettingGroup()); @@ -75,7 +93,8 @@ public InterpreterClient launchDirectly(InterpreterLaunchContext context) throws zConf.getInterpreterDir() + "/" + groupName, localRepoPath, buildEnvFromProperties(context), connectTimeout, connectionPoolSize, name, context.getInterpreterGroupId(), option.isUserImpersonate(), - runner != null ? runner.getPath() : zConf.getInterpreterRemoteRunnerPath()); + runner != null ? runner.getPath() : zConf.getInterpreterRemoteRunnerPath(), + processLaunchObserver); } } diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/recovery/StopInterpreter.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/recovery/StopInterpreter.java index d7b710a1305..d86309e6e7b 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/recovery/StopInterpreter.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/recovery/StopInterpreter.java @@ -23,7 +23,6 @@ import org.apache.zeppelin.interpreter.launcher.InterpreterClient; import org.apache.zeppelin.plugin.PluginManager; import org.apache.zeppelin.storage.ConfigStorage; -import org.apache.zeppelin.util.ReflectionUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -41,15 +40,15 @@ public class StopInterpreter { private static final Logger LOGGER = LoggerFactory.getLogger(StopInterpreter.class); public StopInterpreter(ZeppelinConfiguration zConf) throws IOException { - ConfigStorage storage = ConfigStorage.createConfigStorage(zConf); PluginManager pluginManager = new PluginManager(zConf); + ConfigStorage storage = ConfigStorage.createConfigStorage(zConf, pluginManager); InterpreterSettingManager interpreterSettingManager = new InterpreterSettingManager(zConf, null, null, null, storage, pluginManager); - RecoveryStorage recoveryStorage = - ReflectionUtils.createClazzInstance(zConf.getRecoveryStorageClass(), - new Class[] { ZeppelinConfiguration.class, InterpreterSettingManager.class }, - new Object[] { zConf, interpreterSettingManager }); + RecoveryStorage recoveryStorage = pluginManager.createPluginInstance( + zConf.getRecoveryStorageClass(), + new Class[] { ZeppelinConfiguration.class, InterpreterSettingManager.class }, + new Object[] { zConf, interpreterSettingManager }); LOGGER.info("Using RecoveryStorage: {}", recoveryStorage.getClass().getName()); Map restoredClients = recoveryStorage.restore(); diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcess.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcess.java index 6dd2793adf8..d43d902db0e 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcess.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcess.java @@ -19,13 +19,9 @@ import java.io.IOException; import java.util.Map; -import java.util.regex.Matcher; -import java.util.regex.Pattern; import org.apache.commons.exec.CommandLine; import org.apache.commons.exec.ExecuteException; -import org.apache.hadoop.yarn.util.ConverterUtils; -import org.apache.zeppelin.interpreter.YarnAppMonitor; import org.apache.zeppelin.interpreter.util.ProcessLauncher; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -36,9 +32,8 @@ public class ExecRemoteInterpreterProcess extends RemoteInterpreterManagedProces private static final Logger LOGGER = LoggerFactory.getLogger(ExecRemoteInterpreterProcess.class); - private static final Pattern YARN_APP_PATTER = Pattern.compile("Submitted application (\\w+)"); - private final String interpreterRunner; + private final ProcessLaunchObserver processLaunchObserver; private InterpreterProcessLauncher interpreterProcessLauncher; public ExecRemoteInterpreterProcess( @@ -54,9 +49,50 @@ public ExecRemoteInterpreterProcess( String interpreterGroupId, boolean isUserImpersonated, String intpRunner) { - super(intpEventServerPort, intpEventServerHost, interpreterPortRange, intpDir, localRepoDir, env, connectTimeout, - connectionPoolSize, interpreterSettingName, interpreterGroupId, isUserImpersonated); + this( + intpEventServerPort, + intpEventServerHost, + interpreterPortRange, + intpDir, + localRepoDir, + env, + connectTimeout, + connectionPoolSize, + interpreterSettingName, + interpreterGroupId, + isUserImpersonated, + intpRunner, + ProcessLaunchObserver.NO_OP); + } + + public ExecRemoteInterpreterProcess( + int intpEventServerPort, + String intpEventServerHost, + String interpreterPortRange, + String intpDir, + String localRepoDir, + Map env, + int connectTimeout, + int connectionPoolSize, + String interpreterSettingName, + String interpreterGroupId, + boolean isUserImpersonated, + String intpRunner, + ProcessLaunchObserver processLaunchObserver) { + super( + intpEventServerPort, + intpEventServerHost, + interpreterPortRange, + intpDir, + localRepoDir, + env, + connectTimeout, + connectionPoolSize, + interpreterSettingName, + interpreterGroupId, + isUserImpersonated); this.interpreterRunner = intpRunner; + this.processLaunchObserver = processLaunchObserver; } @Override @@ -87,34 +123,23 @@ public void start(String userName) throws IOException { interpreterProcessLauncher.waitForReady(getConnectTimeout()); if (interpreterProcessLauncher.isLaunchTimeout()) { throw new IOException( - String.format("Interpreter Process creation is time out in %d seconds", getConnectTimeout() / 1000) + "\n" + String.format( + "Interpreter Process creation is time out in %d seconds", + getConnectTimeout() / 1000) + + "\n" + "You can increase timeout threshold via " + "setting zeppelin.interpreter.connect.timeout of this interpreter.\n" + interpreterProcessLauncher.getErrorMessage()); } if (!interpreterProcessLauncher.isRunning()) { - throw new IOException("Fail to launch interpreter process:\n" + interpreterProcessLauncher.getErrorMessage()); - } - - if (isHadoopClientAvailable()) { - String launchOutput = interpreterProcessLauncher.getProcessLaunchOutput(); - Matcher m = YARN_APP_PATTER.matcher(launchOutput); - if (m.find()) { - String appId = m.group(1); - LOGGER.info("Detected yarn app: {}, add it to YarnAppMonitor", appId); - YarnAppMonitor.get().addYarnApp(ConverterUtils.toApplicationId(appId), this); - } + throw new IOException( + "Fail to launch interpreter process:\n" + + interpreterProcessLauncher.getErrorMessage()); } - } - private boolean isHadoopClientAvailable() { - try { - Class.forName("org.apache.hadoop.yarn.conf.YarnConfiguration"); - return true; - } catch (ClassNotFoundException e) { - return false; - } + processLaunchObserver.onProcessLaunch( + interpreterProcessLauncher.getProcessLaunchOutput(), this); } @Override @@ -129,15 +154,20 @@ public void stop() { if (isRunning()) { super.stop(); // wait for a clean shutdown - this.interpreterProcessLauncher.waitForShutdown(RemoteInterpreterServer.DEFAULT_SHUTDOWN_TIMEOUT + 500); + this.interpreterProcessLauncher.waitForShutdown( + RemoteInterpreterServer.DEFAULT_SHUTDOWN_TIMEOUT + 500); // kill process this.interpreterProcessLauncher.stop(); this.interpreterProcessLauncher = null; - LOGGER.info("Remote exec process of interpreter group: {} is terminated", getInterpreterGroupId()); + LOGGER.info( + "Remote exec process of interpreter group: {} is terminated", + getInterpreterGroupId()); } else { // Shutdown connection super.close(); - LOGGER.warn("Try to stop a not running interpreter process of interpreter group: {}", getInterpreterGroupId()); + LOGGER.warn( + "Try to stop a not running interpreter process of interpreter group: {}", + getInterpreterGroupId()); } } @@ -165,7 +195,7 @@ public String getErrorMessage() { private class InterpreterProcessLauncher extends ProcessLauncher { - public InterpreterProcessLauncher(CommandLine commandLine, Map envs) { + InterpreterProcessLauncher(CommandLine commandLine, Map envs) { super(commandLine, envs); } diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/ProcessLaunchObserver.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/ProcessLaunchObserver.java new file mode 100644 index 00000000000..a2c92a80a0a --- /dev/null +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/ProcessLaunchObserver.java @@ -0,0 +1,34 @@ +/* + * 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.zeppelin.interpreter.remote; + +/** + * Observes the output produced while an interpreter process is launched. + * + *

Implementations can use the output to discover an external application and monitor its + * lifecycle. This interface deliberately exposes no cluster-manager-specific types so that + * implementations and their dependencies can live in optional plugins. + */ +@FunctionalInterface +public interface ProcessLaunchObserver { + + ProcessLaunchObserver NO_OP = (launchOutput, interpreterProcess) -> { }; + + void onProcessLaunch( + String launchOutput, RemoteInterpreterManagedProcess interpreterProcess); +} diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/notebook/repo/NotebookRepoSync.java b/zeppelin-server/src/main/java/org/apache/zeppelin/notebook/repo/NotebookRepoSync.java index eb5b1e37e58..ccce415d7b2 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/notebook/repo/NotebookRepoSync.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/notebook/repo/NotebookRepoSync.java @@ -82,7 +82,7 @@ public void init(ZeppelinConfiguration zConf, NoteParser noteParser) throws IOEx for (int i = 0; i < Math.min(storageClassNames.length, getMaxRepoNum()); i++) { NotebookRepo notebookRepo = pluginManager.loadNotebookRepo(storageClassNames[i].trim()); - notebookRepo.init(zConf, noteParser); + initNotebookRepo(notebookRepo, zConf, noteParser); repos.add(notebookRepo); } @@ -90,7 +90,7 @@ public void init(ZeppelinConfiguration zConf, NoteParser noteParser) throws IOEx if (getRepoCount() == 0) { LOGGER.info("No storage could be initialized, using default {} storage", DEFAULT_STORAGE); NotebookRepo defaultNotebookRepo = pluginManager.loadNotebookRepo(DEFAULT_STORAGE); - defaultNotebookRepo.init(zConf, noteParser); + initNotebookRepo(defaultNotebookRepo, zConf, noteParser); repos.add(defaultNotebookRepo); } // sync for anonymous mode on start @@ -103,6 +103,19 @@ public void init(ZeppelinConfiguration zConf, NoteParser noteParser) throws IOEx } } + private void initNotebookRepo( + NotebookRepo notebookRepo, ZeppelinConfiguration zConf, NoteParser noteParser) + throws IOException { + Thread thread = Thread.currentThread(); + ClassLoader previousClassLoader = thread.getContextClassLoader(); + try { + thread.setContextClassLoader(notebookRepo.getClass().getClassLoader()); + notebookRepo.init(zConf, noteParser); + } finally { + thread.setContextClassLoader(previousClassLoader); + } + } + public List getNotebookRepos(AuthenticationInfo subject) { List reposSetting = new ArrayList<>(); NotebookRepoWithSettings repoWithSettings; diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/plugin/PluginManager.java b/zeppelin-server/src/main/java/org/apache/zeppelin/plugin/PluginManager.java index b229d0f0f57..07c22011c8d 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/plugin/PluginManager.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/plugin/PluginManager.java @@ -22,6 +22,7 @@ import org.apache.zeppelin.interpreter.launcher.SparkInterpreterLauncher; import org.apache.zeppelin.interpreter.launcher.StandardInterpreterLauncher; import org.apache.zeppelin.interpreter.recovery.RecoveryStorage; +import org.apache.zeppelin.interpreter.remote.ProcessLaunchObserver; import org.apache.zeppelin.notebook.repo.GitNotebookRepo; import org.apache.zeppelin.notebook.repo.NotebookRepo; import org.apache.zeppelin.notebook.repo.VFSNotebookRepo; @@ -35,9 +36,13 @@ import java.net.URLClassLoader; import java.util.ArrayList; import java.util.Arrays; +import java.util.Collections; import java.util.HashMap; +import java.util.LinkedHashSet; import java.util.List; import java.util.Map; +import java.util.ServiceLoader; +import java.util.Set; import jakarta.inject.Inject; @@ -51,7 +56,8 @@ public class PluginManager { private final String pluginsDir; private final ZeppelinConfiguration zConf; - private Map cachedLaunchers = new HashMap<>(); + private final Map cachedLaunchers = new HashMap<>(); + private final Map pluginClassLoaders = new HashMap<>(); private List builtinLauncherClassNames = Arrays.asList( StandardInterpreterLauncher.class.getName(), @@ -85,10 +91,13 @@ public NotebookRepo loadNotebookRepo(String notebookRepoClassName) throws IOExce } NotebookRepo notebookRepo = null; try { - notebookRepo = (NotebookRepo) (Class.forName(notebookRepoClassName, true, pluginClassLoader)).newInstance(); - } catch (InstantiationException | IllegalAccessException | ClassNotFoundException e) { + notebookRepo = withContextClassLoader(pluginClassLoader, () -> + (NotebookRepo) Class.forName(notebookRepoClassName, true, pluginClassLoader) + .getDeclaredConstructor() + .newInstance()); + } catch (ReflectiveOperationException e) { throw new IOException("Fail to instantiate notebookrepo " + notebookRepoClassName + - " from plugin classpath:" + pluginsDir, e); + " from plugin classpath:" + pluginsDir, e); } return notebookRepo; @@ -112,10 +121,12 @@ public synchronized InterpreterLauncher loadInterpreterLauncher(String launcherP if (builtinLauncherClassNames.contains(launcherClassName) || Boolean.parseBoolean(System.getProperty("zeppelin.isTest", "false"))) { try { - return (InterpreterLauncher) + InterpreterLauncher launcher = (InterpreterLauncher) (Class.forName(launcherClassName)) .getConstructor(ZeppelinConfiguration.class, RecoveryStorage.class) .newInstance(zConf, recoveryStorage); + configureProcessLaunchObservers(launcher); + return launcher; } catch (InstantiationException | IllegalAccessException | ClassNotFoundException | NoSuchMethodException | InvocationTargetException e) { throw new IOException("Fail to instantiate InterpreterLauncher from classpath directly:" @@ -126,37 +137,184 @@ public synchronized InterpreterLauncher loadInterpreterLauncher(String launcherP URLClassLoader pluginClassLoader = getPluginClassLoader(pluginsDir, "Launcher", launcherPlugin); InterpreterLauncher launcher = null; try { - launcher = (InterpreterLauncher) (Class.forName(launcherClassName, true, pluginClassLoader)) - .getConstructor(ZeppelinConfiguration.class, RecoveryStorage.class) - .newInstance(zConf, recoveryStorage); - } catch (InstantiationException | IllegalAccessException | ClassNotFoundException - | NoSuchMethodException | InvocationTargetException e) { + launcher = withContextClassLoader(pluginClassLoader, () -> + (InterpreterLauncher) Class.forName(launcherClassName, true, pluginClassLoader) + .getConstructor(ZeppelinConfiguration.class, RecoveryStorage.class) + .newInstance(zConf, recoveryStorage)); + } catch (ReflectiveOperationException e) { throw new IOException("Fail to instantiate Launcher " + launcherPlugin + " from plugin pluginDir: " + pluginsDir, e); } + configureProcessLaunchObservers(launcher); cachedLaunchers.put(launcherPlugin, launcher); return launcher; } + private void configureProcessLaunchObservers(InterpreterLauncher launcher) throws IOException { + if (launcher instanceof StandardInterpreterLauncher) { + ((StandardInterpreterLauncher) launcher) + .setProcessLaunchObservers(loadServiceProviders(ProcessLaunchObserver.class)); + } + } + private URLClassLoader getPluginClassLoader(String pluginsDir, String pluginType, String pluginName) throws IOException { File pluginFolder = new File(pluginsDir + "/" + pluginType + "/" + pluginName); + return getPluginClassLoader(pluginFolder); + } + + private synchronized URLClassLoader getPluginClassLoader(File pluginFolder) throws IOException { if (!pluginFolder.exists() || pluginFolder.isFile()) { LOGGER.warn("PluginFolder {} doesn't exist or is not a directory", pluginFolder.getAbsolutePath()); return null; } + String pluginFolderPath = pluginFolder.getCanonicalPath(); + if (pluginClassLoaders.containsKey(pluginFolderPath)) { + return pluginClassLoaders.get(pluginFolderPath); + } List urls = new ArrayList<>(); - for (File file : pluginFolder.listFiles()) { - LOGGER.debug("Add file {} to classpath of plugin: {}", file.getAbsolutePath(), pluginName); - urls.add(file.toURI().toURL()); + File[] pluginFiles = pluginFolder.listFiles(); + if (pluginFiles != null) { + for (File file : pluginFiles) { + LOGGER.debug("Add file {} to classpath of plugin: {}", + file.getAbsolutePath(), pluginFolder.getName()); + urls.add(file.toURI().toURL()); + } } if (urls.isEmpty()) { - LOGGER.warn("Can not load plugin {}, because the plugin folder {} is empty.", pluginName , pluginFolder); + LOGGER.warn("Can not load plugin, because the plugin folder {} is empty.", pluginFolder); return null; } - return new URLClassLoader(urls.toArray(new URL[0])); + URLClassLoader classLoader = new URLClassLoader( + urls.toArray(new URL[0]), PluginManager.class.getClassLoader()); + pluginClassLoaders.put(pluginFolderPath, classLoader); + return classLoader; + } + + /** + * Load an extension class from any plugin directory. + * + *

This is used by configurable extension points such as ConfigStorage and RecoveryStorage, + * whose implementation class name is stored in zeppelin-site.xml. The implementation and all + * of its third-party dependencies stay in the plugin classloader.

+ */ + public Class loadPluginClass(String className) throws IOException { + try { + return Class.forName(className, true, PluginManager.class.getClassLoader()); + } catch (ClassNotFoundException e) { + File pluginFolder = findPluginFolder(className); + if (pluginFolder == null) { + throw new IOException("Unable to find plugin class: " + className, e); + } + try { + URLClassLoader classLoader = getPluginClassLoader(pluginFolder); + return withContextClassLoader( + classLoader, () -> Class.forName(className, true, classLoader)); + } catch (ClassNotFoundException pluginError) { + throw new IOException("Unable to load plugin class: " + className, pluginError); + } catch (ReflectiveOperationException pluginError) { + throw new IOException("Unable to initialize plugin class: " + className, pluginError); + } + } + } + + public T createPluginInstance(String className, + Class[] parameterTypes, + Object[] parameters) throws IOException { + try { + Class pluginClass = loadPluginClass(className); + @SuppressWarnings("unchecked") T instance = withContextClassLoader( + pluginClass.getClassLoader(), + () -> (T) pluginClass.getConstructor(parameterTypes).newInstance(parameters)); + return instance; + } catch (ReflectiveOperationException e) { + throw new IOException("Unable to instantiate plugin class: " + className, e); + } + } + + /** Load service providers without adding their dependency jars to the server classpath. */ + public List loadServiceProviders(Class serviceType) throws IOException { + List providers = new ArrayList<>(); + Set providerClassNames = new LinkedHashSet<>(); + for (File pluginFolder : getPluginFolders()) { + URLClassLoader classLoader = getPluginClassLoader(pluginFolder); + if (classLoader == null) { + continue; + } + Thread thread = Thread.currentThread(); + ClassLoader previousClassLoader = thread.getContextClassLoader(); + try { + thread.setContextClassLoader(classLoader); + for (T provider : ServiceLoader.load(serviceType, classLoader)) { + if (provider.getClass().getClassLoader() == classLoader && + providerClassNames.add(provider.getClass().getName())) { + providers.add(provider); + } + } + } finally { + thread.setContextClassLoader(previousClassLoader); + } + } + return providers; + } + + /** Return the isolated classpath containing a configured plugin class. */ + public List getPluginClasspath(String className) throws IOException { + File pluginFolder = findPluginFolder(className); + if (pluginFolder == null) { + return Collections.emptyList(); + } + File[] files = pluginFolder.listFiles(); + if (files == null) { + return Collections.emptyList(); + } + return Arrays.asList(files); + } + + private File findPluginFolder(String className) throws IOException { + String classResource = className.replace('.', '/') + ".class"; + for (File pluginFolder : getPluginFolders()) { + URLClassLoader classLoader = getPluginClassLoader(pluginFolder); + if (classLoader != null && classLoader.findResource(classResource) != null) { + return pluginFolder; + } + } + return null; + } + + private List getPluginFolders() { + File root = new File(pluginsDir); + File[] pluginTypes = root.listFiles(File::isDirectory); + if (pluginTypes == null) { + return Collections.emptyList(); + } + List pluginFolders = new ArrayList<>(); + for (File pluginType : pluginTypes) { + File[] folders = pluginType.listFiles(File::isDirectory); + if (folders != null) { + pluginFolders.addAll(Arrays.asList(folders)); + } + } + return pluginFolders; + } + + @FunctionalInterface + private interface ReflectiveAction { + T run() throws ReflectiveOperationException; + } + + private static T withContextClassLoader( + ClassLoader classLoader, ReflectiveAction action) throws ReflectiveOperationException { + Thread thread = Thread.currentThread(); + ClassLoader previousClassLoader = thread.getContextClassLoader(); + try { + thread.setContextClassLoader(classLoader); + return action.run(); + } finally { + thread.setContextClassLoader(previousClassLoader); + } } } diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/realm/ExternalLoginRealm.java b/zeppelin-server/src/main/java/org/apache/zeppelin/realm/ExternalLoginRealm.java new file mode 100644 index 00000000000..350e1a5c24e --- /dev/null +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/realm/ExternalLoginRealm.java @@ -0,0 +1,49 @@ +/* + * 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.zeppelin.realm; + +import java.util.Map; + +import jakarta.ws.rs.core.Cookie; +import org.apache.shiro.authc.AuthenticationException; +import org.apache.shiro.authc.AuthenticationToken; + +/** + * Contract used by the server login endpoint to interact with optional SSO realms without + * depending on their implementation classes. + */ +public interface ExternalLoginRealm { + + AuthenticationToken getLoginAuthenticationToken(Map cookies) + throws AuthenticationException; + + String getLoginPrincipal(AuthenticationToken token) throws AuthenticationException; + + boolean shouldRedirectOnMissingToken(); + + int getLoginPriority(); + + String getProviderUrl(); + + String getRedirectParam(); + + String getLogin(); + + String getLogout(); + + Boolean getLogoutAPI(); +} diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/realm/GroupResolver.java b/zeppelin-server/src/main/java/org/apache/zeppelin/realm/GroupResolver.java new file mode 100644 index 00000000000..db70f23abc3 --- /dev/null +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/realm/GroupResolver.java @@ -0,0 +1,26 @@ +/* + * 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.zeppelin.realm; + +import java.io.IOException; +import java.util.Set; + +/** Resolves the external groups associated with a user. */ +public interface GroupResolver { + + Set resolve(String principal) throws IOException; +} diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/realm/LdapRealm.java b/zeppelin-server/src/main/java/org/apache/zeppelin/realm/LdapRealm.java index be8a0f0c68d..c01c72d93f1 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/realm/LdapRealm.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/realm/LdapRealm.java @@ -44,9 +44,6 @@ import javax.naming.ldap.LdapContext; import javax.naming.ldap.LdapName; import javax.naming.ldap.PagedResultsControl; -import org.apache.hadoop.conf.Configuration; -import org.apache.hadoop.security.alias.CredentialProvider; -import org.apache.hadoop.security.alias.CredentialProviderFactory; import org.apache.shiro.SecurityUtils; import org.apache.shiro.ShiroException; import org.apache.shiro.authc.AuthenticationInfo; @@ -184,6 +181,8 @@ public class LdapRealm extends DefaultLdapRealm { private String hadoopSecurityCredentialPath; private static final String KEYSTORE_PASS = "ldapRealm.systemPassword"; + private static final String HADOOP_SECRET_RESOLVER = + "org.apache.zeppelin.realm.hadoop.HadoopCredentialProviderSecretResolver"; private boolean authorizationEnabled; @@ -228,15 +227,13 @@ static String getSystemPassword(String hadoopSecurityCredentialPath, String keystorePass) { String password = ""; try { - Configuration configuration = new Configuration(); - configuration.set(CredentialProviderFactory.CREDENTIAL_PROVIDER_PATH, - hadoopSecurityCredentialPath); - CredentialProvider provider = CredentialProviderFactory.getProviders(configuration).get(0); - CredentialProvider.CredentialEntry credEntry = provider.getCredentialEntry(keystorePass); - if (credEntry != null) { - password = new String(credEntry.getCredential()); + SecretResolver resolver = + SecurityProviderLoader.load(HADOOP_SECRET_RESOLVER, SecretResolver.class); + char[] credential = resolver.resolve(hadoopSecurityCredentialPath, keystorePass); + if (credential != null) { + password = new String(credential); } - } catch (IOException e) { + } catch (IOException | ReflectiveOperationException e) { throw new ShiroException("Error from getting credential entry from keystore", e); } if (org.apache.commons.lang3.StringUtils.isEmpty(password)) { diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/realm/SecretResolver.java b/zeppelin-server/src/main/java/org/apache/zeppelin/realm/SecretResolver.java new file mode 100644 index 00000000000..7f1c34edf64 --- /dev/null +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/realm/SecretResolver.java @@ -0,0 +1,25 @@ +/* + * 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.zeppelin.realm; + +import java.io.IOException; + +/** Resolves a named secret from an external provider. */ +public interface SecretResolver { + + char[] resolve(String providerPath, String alias) throws IOException; +} diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/realm/SecurityProviderLoader.java b/zeppelin-server/src/main/java/org/apache/zeppelin/realm/SecurityProviderLoader.java new file mode 100644 index 00000000000..24b37d8c13c --- /dev/null +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/realm/SecurityProviderLoader.java @@ -0,0 +1,34 @@ +/* + * 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.zeppelin.realm; + +/** Loads optional security providers from the class loader that created the Shiro environment. */ +public final class SecurityProviderLoader { + + private SecurityProviderLoader() { + } + + public static T load(String className, Class providerType) + throws ReflectiveOperationException { + ClassLoader classLoader = Thread.currentThread().getContextClassLoader(); + if (classLoader == null) { + classLoader = SecurityProviderLoader.class.getClassLoader(); + } + Class providerClass = Class.forName(className, true, classLoader); + return providerType.cast(providerClass.getDeclaredConstructor().newInstance()); + } +} diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/realm/ZeppelinRoleProvider.java b/zeppelin-server/src/main/java/org/apache/zeppelin/realm/ZeppelinRoleProvider.java new file mode 100644 index 00000000000..c7c76701385 --- /dev/null +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/realm/ZeppelinRoleProvider.java @@ -0,0 +1,25 @@ +/* + * 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.zeppelin.realm; + +import java.util.Set; + +/** Supplies the Zeppelin roles associated with an authenticated principal. */ +public interface ZeppelinRoleProvider { + + Set mapGroupPrincipals(String principal); +} diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/realm/jwt/KnoxJwtRealm.java b/zeppelin-server/src/main/java/org/apache/zeppelin/realm/jwt/KnoxJwtRealm.java index 0a1d8decc9c..98419406dd5 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/realm/jwt/KnoxJwtRealm.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/realm/jwt/KnoxJwtRealm.java @@ -16,48 +16,55 @@ */ package org.apache.zeppelin.realm.jwt; -import java.nio.charset.Charset; -import java.nio.charset.StandardCharsets; -import java.util.Date; -import org.apache.commons.io.FileUtils; -import org.apache.hadoop.conf.Configuration; -import org.apache.hadoop.security.Groups; -import org.apache.shiro.authc.AuthenticationInfo; -import org.apache.shiro.authc.AuthenticationToken; -import org.apache.shiro.authc.SimpleAccount; -import org.apache.shiro.authz.AuthorizationInfo; -import org.apache.shiro.authz.SimpleAuthorizationInfo; -import org.apache.shiro.realm.AuthorizingRealm; -import org.apache.shiro.subject.PrincipalCollection; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - import java.io.ByteArrayInputStream; import java.io.File; import java.io.IOException; import java.io.UnsupportedEncodingException; +import java.nio.charset.Charset; +import java.nio.charset.StandardCharsets; import java.security.PublicKey; import java.security.cert.CertificateException; import java.security.cert.CertificateFactory; import java.security.cert.X509Certificate; import java.security.interfaces.RSAPublicKey; import java.text.ParseException; -import java.util.HashSet; -import java.util.List; +import java.util.Collections; +import java.util.Date; +import java.util.Map; import java.util.Set; import jakarta.servlet.ServletException; +import jakarta.ws.rs.core.Cookie; import com.nimbusds.jose.JWSObject; import com.nimbusds.jose.JWSVerifier; import com.nimbusds.jose.crypto.RSASSAVerifier; import com.nimbusds.jwt.SignedJWT; +import org.apache.commons.io.FileUtils; +import org.apache.shiro.ShiroException; +import org.apache.shiro.authc.AuthenticationException; +import org.apache.shiro.authc.AuthenticationInfo; +import org.apache.shiro.authc.AuthenticationToken; +import org.apache.shiro.authc.SimpleAccount; +import org.apache.shiro.authz.AuthorizationInfo; +import org.apache.shiro.authz.SimpleAuthorizationInfo; +import org.apache.shiro.realm.AuthorizingRealm; +import org.apache.shiro.subject.PrincipalCollection; +import org.apache.zeppelin.realm.ExternalLoginRealm; +import org.apache.zeppelin.realm.GroupResolver; +import org.apache.zeppelin.realm.SecurityProviderLoader; +import org.apache.zeppelin.realm.ZeppelinRoleProvider; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * Created for org.apache.zeppelin.server. */ -public class KnoxJwtRealm extends AuthorizingRealm { +public class KnoxJwtRealm extends AuthorizingRealm + implements ExternalLoginRealm, ZeppelinRoleProvider { private static final Logger LOGGER = LoggerFactory.getLogger(KnoxJwtRealm.class); + private static final String HADOOP_GROUP_RESOLVER = + "org.apache.zeppelin.realm.hadoop.HadoopGroupResolver"; private String providerUrl; private String redirectParam; @@ -66,21 +73,19 @@ public class KnoxJwtRealm extends AuthorizingRealm { private String login; private String logout; private Boolean logoutAPI; + private String groupResolverClass = HADOOP_GROUP_RESOLVER; - /** - * Hadoop Groups implementation. - */ - private Groups hadoopGroups; + private GroupResolver groupResolver = principal -> Collections.emptySet(); @Override protected void onInit() { super.onInit(); try { - Configuration hadoopConfig = new Configuration(); - hadoopGroups = new Groups(hadoopConfig); + groupResolver = SecurityProviderLoader.load(groupResolverClass, GroupResolver.class); } catch (final Exception e) { - LOGGER.error("Exception in onInit", e); + throw new ShiroException( + "Unable to load the Knox group resolver: " + groupResolverClass, e); } } @@ -215,22 +220,17 @@ protected AuthorizationInfo doGetAuthorizationInfo(PrincipalCollection principal } /** - * Query the Hadoop implementation of {@link Groups} to retrieve groups for provided user. + * Query the configured resolver to retrieve groups for the provided user. */ public Set mapGroupPrincipals(final String mappedPrincipalName) { - /* return the groups as seen by Hadoop */ - Set groups; try { - final List groupList = hadoopGroups - .getGroups(mappedPrincipalName); + Set groups = groupResolver.resolve(mappedPrincipalName); if (LOGGER.isDebugEnabled()) { LOGGER.debug(String.format("group found %s, %s", - mappedPrincipalName, groupList.toString())); + mappedPrincipalName, groups.toString())); } - - groups = new HashSet<>(groupList); - + return groups; } catch (final IOException e) { if (e.toString().contains("No groups found for user")) { /* no groups found move on */ @@ -240,9 +240,49 @@ public Set mapGroupPrincipals(final String mappedPrincipalName) { /* Log the error and return empty group */ LOGGER.info(String.format("errorGettingUserGroups for %s", mappedPrincipalName)); } - groups = new HashSet<>(); + return Collections.emptySet(); + } + } + + void setGroupResolver(GroupResolver groupResolver) { + this.groupResolver = groupResolver; + } + + @Override + public AuthenticationToken getLoginAuthenticationToken( + Map cookies) { + Cookie cookie = cookies.get(cookieName); + if (cookie == null || cookie.getValue() == null) { + return null; + } + return new JWTAuthenticationToken(null, cookie.getValue()); + } + + @Override + public String getLoginPrincipal(AuthenticationToken token) throws AuthenticationException { + try { + return getName((JWTAuthenticationToken) token); + } catch (ParseException e) { + throw new AuthenticationException("Unable to parse the Knox JWT", e); } - return groups; + } + + @Override + public boolean shouldRedirectOnMissingToken() { + return true; + } + + @Override + public int getLoginPriority() { + return 100; + } + + public String getGroupResolverClass() { + return groupResolverClass; + } + + public void setGroupResolverClass(String groupResolverClass) { + this.groupResolverClass = groupResolverClass; } public String getProviderUrl() { diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/rest/LoginRestApi.java b/zeppelin-server/src/main/java/org/apache/zeppelin/rest/LoginRestApi.java index d8b8c93b93e..d92a9177c8e 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/rest/LoginRestApi.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/rest/LoginRestApi.java @@ -16,7 +16,6 @@ */ package org.apache.zeppelin.rest; -import java.text.ParseException; import java.util.Collection; import java.util.HashMap; import java.util.Map; @@ -29,7 +28,6 @@ import jakarta.ws.rs.Path; import jakarta.ws.rs.Produces; import jakarta.ws.rs.core.Context; -import jakarta.ws.rs.core.Cookie; import jakarta.ws.rs.core.HttpHeaders; import jakarta.ws.rs.core.Response; import jakarta.ws.rs.core.Response.Status; @@ -43,10 +41,7 @@ import org.apache.zeppelin.annotation.ZeppelinApi; import org.apache.zeppelin.conf.ZeppelinConfiguration; import org.apache.zeppelin.notebook.AuthorizationService; -import org.apache.zeppelin.realm.jwt.JWTAuthenticationToken; -import org.apache.zeppelin.realm.jwt.KnoxJwtRealm; -import org.apache.zeppelin.realm.kerberos.KerberosRealm; -import org.apache.zeppelin.realm.kerberos.KerberosToken; +import org.apache.zeppelin.realm.ExternalLoginRealm; import org.apache.zeppelin.server.JsonResponse; import org.apache.zeppelin.service.AuthenticationService; import org.apache.zeppelin.ticket.TicketContainer; @@ -78,108 +73,57 @@ public LoginRestApi(ZeppelinConfiguration zConf, @ZeppelinApi public Response getLogin(@Context HttpHeaders headers) { JsonResponse> response = null; - if (isKnoxSSOEnabled()) { - KnoxJwtRealm knoxJwtRealm = getJTWRealm(); - Cookie cookie = headers.getCookies().get(knoxJwtRealm.getCookieName()); - if (cookie != null && cookie.getValue() != null) { - Subject currentUser = SecurityUtils.getSubject(); - JWTAuthenticationToken token = new JWTAuthenticationToken(null, cookie.getValue()); - try { - String name = knoxJwtRealm.getName(token); + ExternalLoginRealm externalLoginRealm = getExternalLoginRealm(); + if (externalLoginRealm != null) { + try { + AuthenticationToken token = + externalLoginRealm.getLoginAuthenticationToken(headers.getCookies()); + if (token != null) { + String name = externalLoginRealm.getLoginPrincipal(token); + Subject currentUser = SecurityUtils.getSubject(); if (!currentUser.isAuthenticated() || !currentUser.getPrincipal().equals(name)) { response = proceedToLogin(currentUser, token); } - } catch (ParseException e) { - LOGGER.error("ParseException in LoginRestApi: ", e); } + } catch (AuthenticationException e) { + LOGGER.error("Error while processing an external login token", e); } - if (response == null) { + + if (response == null && externalLoginRealm.shouldRedirectOnMissingToken()) { Map data = new HashMap<>(); data.put("redirectURL", - constructUrl(knoxJwtRealm.getProviderUrl(), knoxJwtRealm.getRedirectParam(), - knoxJwtRealm.getLogin())); + constructUrl(externalLoginRealm.getProviderUrl(), + externalLoginRealm.getRedirectParam(), externalLoginRealm.getLogin())); response = new JsonResponse<>(Status.OK, "", data); } - return response.build(); - } - - KerberosRealm kerberosRealm = getKerberosRealm(); - if (null != kerberosRealm) { - try { - Map cookies = headers.getCookies(); - KerberosToken kerberosToken = KerberosRealm.getKerberosTokenFromCookies(cookies); - if (null != kerberosToken) { - Subject currentUser = SecurityUtils.getSubject(); - String name = (String) kerberosToken.getPrincipal(); - if (!currentUser.isAuthenticated() || !currentUser.getPrincipal().equals(name)) { - response = proceedToLogin(currentUser, kerberosToken); - } - } - if (null == response) { - LOGGER.warn("No Kerberos token received"); - response = new JsonResponse<>(Status.UNAUTHORIZED, "", null); - } - return response.build(); - } catch (AuthenticationException e){ - LOGGER.error("Error in Login", e); + if (response == null) { + LOGGER.warn("No external authentication token received"); + response = new JsonResponse<>(Status.UNAUTHORIZED, "", null); } + return response.build(); } return new JsonResponse<>(Status.METHOD_NOT_ALLOWED).build(); } - private KerberosRealm getKerberosRealm() { - Collection realmsList = authenticationService.getRealmsList(); - if (realmsList != null) { - for (Realm realm : realmsList) { - String name = realm.getClass().getName(); - - LOGGER.debug("RealmClass.getName: {}", name); - - if (name.equals("org.apache.zeppelin.realm.kerberos.KerberosRealm")) { - return (KerberosRealm) realm; - } - } - } - return null; - } - - private KnoxJwtRealm getJTWRealm() { + private ExternalLoginRealm getExternalLoginRealm() { + ExternalLoginRealm selectedRealm = null; Collection realmsList = authenticationService.getRealmsList(); if (realmsList != null) { for (Realm realm : realmsList) { - if (realm instanceof KnoxJwtRealm) { - return (KnoxJwtRealm) realm; + LOGGER.debug("RealmClass.getName: {}", realm.getClass().getName()); + if (realm instanceof ExternalLoginRealm + && (selectedRealm == null + || ((ExternalLoginRealm) realm).getLoginPriority() + > selectedRealm.getLoginPriority())) { + selectedRealm = (ExternalLoginRealm) realm; } } } - return null; + return selectedRealm; } - private boolean isKnoxSSOEnabled() { - Collection realmsList = authenticationService.getRealmsList(); - if (realmsList != null) { - for (Realm realm : realmsList) { - if (realm instanceof KnoxJwtRealm) { - return true; - } - } - } - return false; - } - - private boolean isKerberosRealmEnabled() { - Collection realmsList = authenticationService.getRealmsList(); - if (realmsList != null) { - for (Realm realm : realmsList) { - if (realm instanceof KerberosRealm) { - return true; - } - } - } - return false; - } - - private JsonResponse> proceedToLogin(Subject currentUser, AuthenticationToken token) { + private JsonResponse> proceedToLogin( + Subject currentUser, AuthenticationToken token) { JsonResponse> response = null; try { logoutCurrentUser(); @@ -260,18 +204,12 @@ public Response logout() { status = Status.FORBIDDEN; data.put("clearAuthorizationHeader", "false"); } - if (isKnoxSSOEnabled()) { - KnoxJwtRealm knoxJwtRealm = getJTWRealm(); - data.put("redirectURL", - constructUrl(knoxJwtRealm.getProviderUrl(), knoxJwtRealm.getRedirectParam(), - knoxJwtRealm.getLogout())); - data.put("isLogoutAPI", knoxJwtRealm.getLogoutAPI().toString()); - } else if (isKerberosRealmEnabled()) { - KerberosRealm kerberosRealm = getKerberosRealm(); + ExternalLoginRealm externalLoginRealm = getExternalLoginRealm(); + if (externalLoginRealm != null) { data.put("redirectURL", - constructUrl(kerberosRealm.getProviderUrl(), kerberosRealm.getRedirectParam(), - kerberosRealm.getLogout())); - data.put("isLogoutAPI", kerberosRealm.getLogoutAPI().toString()); + constructUrl(externalLoginRealm.getProviderUrl(), externalLoginRealm.getRedirectParam(), + externalLoginRealm.getLogout())); + data.put("isLogoutAPI", externalLoginRealm.getLogoutAPI().toString()); } JsonResponse> response = new JsonResponse<>(status, "", data); LOGGER.info(response.toString()); diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/rest/message/ParagraphJobStatus.java b/zeppelin-server/src/main/java/org/apache/zeppelin/rest/message/ParagraphJobStatus.java index 3d8d4472834..e43d996f4a5 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/rest/message/ParagraphJobStatus.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/rest/message/ParagraphJobStatus.java @@ -17,7 +17,7 @@ package org.apache.zeppelin.rest.message; -import org.apache.commons.lang.StringUtils; +import org.apache.commons.lang3.StringUtils; import org.apache.zeppelin.notebook.Paragraph; import org.apache.zeppelin.scheduler.Job; diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/server/ZeppelinServer.java b/zeppelin-server/src/main/java/org/apache/zeppelin/server/ZeppelinServer.java index b3f78816aec..6d6e5bdde97 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/server/ZeppelinServer.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/server/ZeppelinServer.java @@ -42,14 +42,22 @@ import java.io.IOException; import java.lang.management.ManagementFactory; import java.net.MalformedURLException; +import java.nio.charset.StandardCharsets; import java.nio.file.Files; import java.security.GeneralSecurityException; +import java.util.ArrayList; import java.util.Base64; import java.util.HashSet; +import java.util.LinkedHashMap; +import java.util.LinkedHashSet; import java.util.List; +import java.util.Map; import java.util.Optional; +import java.util.Set; import java.util.concurrent.atomic.AtomicBoolean; import java.util.EnumSet; +import java.util.regex.Matcher; +import java.util.regex.Pattern; import jakarta.inject.Singleton; import javax.management.remote.JMXServiceURL; import jakarta.servlet.DispatcherType; @@ -131,6 +139,15 @@ public class ZeppelinServer implements AutoCloseable { private static final Logger LOGGER = LoggerFactory.getLogger(ZeppelinServer.class); private static final String NON_DEFAULT_NEW_UI_WEB_APP_CONTEXT_PATH = "/new"; private static final String NON_DEFAULT_CLASSIC_UI_WEB_APP_CONTEXT_PATH = "/classic"; + private static final String HADOOP_GROUP_RESOLVER = + "org.apache.zeppelin.realm.hadoop.HadoopGroupResolver"; + private static final String HADOOP_SECRET_RESOLVER = + "org.apache.zeppelin.realm.hadoop.HadoopCredentialProviderSecretResolver"; + private static final Pattern SHIRO_CLASS_ASSIGNMENT = Pattern.compile( + "(?m)^\\s*([A-Za-z_$][\\w$]*)\\s*=\\s*" + + "([A-Za-z_$][\\w$]*(?:\\.[A-Za-z_$][\\w$]*)+)\\s*$"); + private static final Pattern SHIRO_SECTION_HEADER = + Pattern.compile("^\\s*\\[([^]]+)]\\s*$"); public static final String DEFAULT_SERVICE_LOCATOR_NAME = "shared-locator"; private final AtomicBoolean duringShutdown = new AtomicBoolean(false); @@ -138,6 +155,7 @@ public class ZeppelinServer implements AutoCloseable { private final Optional promMetricRegistry; private final Server jettyWebServer; private final ServiceLocator sharedServiceLocator; + private final PluginManager pluginManager; private final ConfigStorage storage; public ZeppelinServer(ZeppelinConfiguration zConf) throws IOException { @@ -154,7 +172,8 @@ public ZeppelinServer(ZeppelinConfiguration zConf, String serviceLocatorName) th } jettyWebServer = setupJettyServer(); sharedServiceLocator = ServiceLocatorFactory.getInstance().create(serviceLocatorName); - storage = ConfigStorage.createConfigStorage(zConf); + pluginManager = new PluginManager(zConf); + storage = ConfigStorage.createConfigStorage(zConf, pluginManager); } public void startZeppelin() { @@ -177,7 +196,7 @@ public void startZeppelin() { @Override protected void configure() { bind(storage).to(ConfigStorage.class); - bindAsContract(PluginManager.class).in(Singleton.class); + bind(pluginManager).to(PluginManager.class); bind(GsonNoteParser.class).to(NoteParser.class).in(Singleton.class); bindAsContract(InterpreterFactory.class).in(Singleton.class); bindAsContract(NotebookRepoSync.class).to(NotebookRepo.class).in(Singleton.class); @@ -562,6 +581,7 @@ private void setupRestApiContextHandler(WebAppContext webapp) { String shiroIniPath = zConf.getShiroPath(); if (!StringUtils.isBlank(shiroIniPath)) { + configureShiroPluginClasspath(webapp, shiroIniPath); webapp.setInitParameter("shiroConfigLocations", new File(shiroIniPath).toURI().toString()); webapp .addFilter(ShiroFilter.class, "/api/*", EnumSet.allOf(DispatcherType.class)) @@ -570,6 +590,121 @@ private void setupRestApiContextHandler(WebAppContext webapp) { } } + private void configureShiroPluginClasspath(WebAppContext webapp, String shiroIniPath) { + try { + String shiroConfig = new String( + Files.readAllBytes(new File(shiroIniPath).toPath()), StandardCharsets.UTF_8); + String activeConfig = String.join("\n", shiroConfig.lines() + .filter(line -> !line.trim().startsWith("#") && !line.trim().startsWith(";")) + .toArray(String[]::new)); + Set pluginClasspath = new LinkedHashSet<>(); + + for (Map.Entry assignment + : findShiroClassAssignments(activeConfig).entrySet()) { + String beanName = assignment.getKey(); + String className = assignment.getValue(); + addPluginClasspathIfRequired(className, pluginClasspath); + if (className.equals("org.apache.zeppelin.realm.jwt.KnoxJwtRealm")) { + String groupResolverClass = findNonBlankAssignment( + activeConfig, beanName + ".groupResolverClass"); + if (groupResolverClass == null) { + addRequiredPluginClasspath(HADOOP_GROUP_RESOLVER, pluginClasspath); + } else { + addPluginClasspathIfRequired(groupResolverClass, pluginClasspath); + } + } + } + if (hasNonBlankAssignment(activeConfig, ".hadoopSecurityCredentialPath")) { + addRequiredPluginClasspath(HADOOP_SECRET_RESOLVER, pluginClasspath); + } + + if (!pluginClasspath.isEmpty()) { + configurePluginClassLoading(webapp, pluginClasspath); + LOGGER.info("Added {} optional security plugin files to the web application classpath", + pluginClasspath.size()); + } + } catch (IOException e) { + throw new IllegalStateException( + "Unable to configure optional security plugins from " + shiroIniPath, e); + } + } + + static Map findShiroClassAssignments(String config) { + Map assignments = new LinkedHashMap<>(); + Matcher matcher = SHIRO_CLASS_ASSIGNMENT.matcher(findShiroMainSection(config)); + while (matcher.find()) { + assignments.put(matcher.group(1), matcher.group(2)); + } + return assignments; + } + + static String findNonBlankAssignment(String config, String propertyName) { + for (String line : findShiroMainSection(config).split("\\R")) { + int separator = line.indexOf('='); + if (separator > 0 && line.substring(0, separator).trim().equals(propertyName)) { + String value = line.substring(separator + 1).trim(); + return StringUtils.isBlank(value) ? null : value; + } + } + return null; + } + + private static String findShiroMainSection(String config) { + StringBuilder mainSection = new StringBuilder(); + boolean inMainSection = false; + for (String line : config.split("\\R")) { + Matcher sectionMatcher = SHIRO_SECTION_HEADER.matcher(line); + if (sectionMatcher.matches()) { + inMainSection = "main".equalsIgnoreCase(sectionMatcher.group(1).trim()); + } else if (inMainSection) { + mainSection.append(line).append('\n'); + } + } + return mainSection.toString(); + } + + static void configurePluginClassLoading(WebAppContext webapp, Set pluginClasspath) + throws IOException { + // Security plugins are loaded by Jetty's child-first WebAppClassLoader. Keep SLF4J on the + // server classloader so a plugin cannot pair its own API with the server's logger binding. + webapp.getSystemClassMatcher().add("org.slf4j."); + List paths = new ArrayList<>(); + for (File file : pluginClasspath) { + paths.add(file.getAbsolutePath()); + } + webapp.setExtraClasspath(String.join(",", paths)); + } + + private void addPluginClasspathIfRequired(String className, Set pluginClasspath) + throws IOException { + try { + Class.forName(className, false, Thread.currentThread().getContextClassLoader()); + } catch (ClassNotFoundException e) { + addRequiredPluginClasspath(className, pluginClasspath); + } + } + + private void addRequiredPluginClasspath(String className, Set pluginClasspath) + throws IOException { + List classpath = pluginManager.getPluginClasspath(className); + if (classpath.isEmpty()) { + throw new IOException("Configured security extension is not installed: " + className); + } + pluginClasspath.addAll(classpath); + } + + private boolean hasNonBlankAssignment(String config, String propertySuffix) { + for (String line : findShiroMainSection(config).split("\\R")) { + int separator = line.indexOf('='); + if (separator > 0 + && line.substring(0, separator).trim().endsWith(propertySuffix) + && StringUtils.isNotBlank(line.substring(separator + 1))) { + return true; + } + } + return false; + } + private void setupPrometheusContextHandler(WebAppContext webapp) { if (promMetricRegistry.isPresent()) { webapp.addServlet(new ServletHolder(new PrometheusServlet(promMetricRegistry.get())), "/metrics"); diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/service/ShiroAuthenticationService.java b/zeppelin-server/src/main/java/org/apache/zeppelin/service/ShiroAuthenticationService.java index 21219c6e2e5..5c6bd9752dd 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/service/ShiroAuthenticationService.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/service/ShiroAuthenticationService.java @@ -54,7 +54,7 @@ import org.apache.zeppelin.conf.ZeppelinConfiguration; import org.apache.zeppelin.realm.ActiveDirectoryGroupRealm; import org.apache.zeppelin.realm.LdapRealm; -import org.apache.zeppelin.realm.jwt.KnoxJwtRealm; +import org.apache.zeppelin.realm.ZeppelinRoleProvider; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -254,8 +254,8 @@ public Set getAssociatedRoles() { } else if (ACTIVE_DIRECTORY_GROUP_REALM.equals(name)) { allRoles = ((ActiveDirectoryGroupRealm) realm).getListRoles(); break; - } else if (realm instanceof KnoxJwtRealm) { - roles = ((KnoxJwtRealm) realm).mapGroupPrincipals(getPrincipal()); + } else if (realm instanceof ZeppelinRoleProvider) { + roles = ((ZeppelinRoleProvider) realm).mapGroupPrincipals(getPrincipal()); break; } } diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/storage/ConfigStorage.java b/zeppelin-server/src/main/java/org/apache/zeppelin/storage/ConfigStorage.java index 6c9692dd8ff..48d0478d149 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/storage/ConfigStorage.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/storage/ConfigStorage.java @@ -24,7 +24,7 @@ import org.apache.zeppelin.interpreter.InterpreterInfoSaving; import org.apache.zeppelin.interpreter.InterpreterSetting; import org.apache.zeppelin.notebook.NotebookAuthorizationInfoSaving; -import org.apache.zeppelin.util.ReflectionUtils; +import org.apache.zeppelin.plugin.PluginManager; import java.io.IOException; @@ -45,9 +45,14 @@ public abstract class ConfigStorage { protected ZeppelinConfiguration zConf; public static ConfigStorage createConfigStorage(ZeppelinConfiguration zConf) throws IOException { + return createConfigStorage(zConf, new PluginManager(zConf)); + } + + public static ConfigStorage createConfigStorage(ZeppelinConfiguration zConf, + PluginManager pluginManager) throws IOException { String configStorageClass = zConf.getString(ZeppelinConfiguration.ConfVars.ZEPPELIN_CONFIG_STORAGE_CLASS); - return ReflectionUtils.createClazzInstance(configStorageClass, + return pluginManager.createPluginInstance(configStorageClass, new Class[] {ZeppelinConfiguration.class}, new Object[] {zConf}); } diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/launcher/InterpreterLauncherTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/launcher/InterpreterLauncherTest.java index e736608fdfc..12e85ae4706 100644 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/launcher/InterpreterLauncherTest.java +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/launcher/InterpreterLauncherTest.java @@ -18,7 +18,13 @@ package org.apache.zeppelin.interpreter.launcher; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertSame; +import java.io.IOException; +import java.net.URL; +import java.net.URLClassLoader; + +import org.apache.zeppelin.conf.ZeppelinConfiguration; import org.junit.jupiter.api.Test; public class InterpreterLauncherTest { @@ -28,4 +34,28 @@ public void testEscapeSpecialCharacters() { String cmd = "{}."; assertEquals("\\{\\}\\.", InterpreterLauncher.escapeSpecialCharacter(cmd)); } + + @Test + void launchUsesPluginClassLoaderAsContextAndRestoresThePreviousOne() throws Exception { + InterpreterLauncher launcher = new InterpreterLauncher(ZeppelinConfiguration.load(), null) { + @Override + public InterpreterClient launchDirectly(InterpreterLaunchContext context) + throws IOException { + assertSame(getClass().getClassLoader(), + Thread.currentThread().getContextClassLoader()); + return null; + } + }; + Thread thread = Thread.currentThread(); + ClassLoader previousClassLoader = thread.getContextClassLoader(); + try (URLClassLoader emptyClassLoader = new URLClassLoader(new URL[0], null)) { + thread.setContextClassLoader(emptyClassLoader); + + launcher.launch(null); + + assertSame(emptyClassLoader, thread.getContextClassLoader()); + } finally { + thread.setContextClassLoader(previousClassLoader); + } + } } diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/recovery/FileSystemRecoveryStorageTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/recovery/FileSystemRecoveryStorageTest.java deleted file mode 100644 index 2b321908c53..00000000000 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/recovery/FileSystemRecoveryStorageTest.java +++ /dev/null @@ -1,118 +0,0 @@ -/* - * 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.zeppelin.interpreter.recovery; - -import org.apache.commons.io.FileUtils; -import org.apache.zeppelin.conf.ZeppelinConfiguration; -import org.apache.zeppelin.interpreter.AbstractInterpreterTest; -import org.apache.zeppelin.interpreter.Interpreter; -import org.apache.zeppelin.interpreter.InterpreterContext; -import org.apache.zeppelin.interpreter.InterpreterException; -import org.apache.zeppelin.interpreter.InterpreterOption; -import org.apache.zeppelin.interpreter.InterpreterSetting; -import org.apache.zeppelin.interpreter.remote.RemoteInterpreter; -import org.apache.zeppelin.user.AuthenticationInfo; -import org.junit.jupiter.api.AfterEach; -import org.junit.jupiter.api.BeforeEach; -import org.junit.jupiter.api.Test; - -import java.io.File; -import java.io.IOException; -import java.nio.file.Files; - -import static org.junit.jupiter.api.Assertions.assertEquals; - -class FileSystemRecoveryStorageTest extends AbstractInterpreterTest { - - private File recoveryDir = null; - private String note1Id; - private String note2Id; - - @Override - @BeforeEach - public void setUp() throws Exception { - System.setProperty(ZeppelinConfiguration.ConfVars.ZEPPELIN_RECOVERY_STORAGE_CLASS.getVarName(), - FileSystemRecoveryStorage.class.getName()); - recoveryDir = Files.createTempDirectory("recoveryDir").toFile(); - System.setProperty(ZeppelinConfiguration.ConfVars.ZEPPELIN_RECOVERY_DIR.getVarName(), recoveryDir.getAbsolutePath()); - super.setUp(); - - note1Id = notebook.createNote("/note_1", AuthenticationInfo.ANONYMOUS); - note2Id = notebook.createNote("/note_2", AuthenticationInfo.ANONYMOUS); - } - - @Override - @AfterEach - public void tearDown() throws Exception { - super.tearDown(); - FileUtils.deleteDirectory(recoveryDir); - System.clearProperty(ZeppelinConfiguration.ConfVars.ZEPPELIN_RECOVERY_STORAGE_CLASS.getVarName()); - } - - @Test - void testSingleInterpreterProcess() throws InterpreterException, IOException { - InterpreterSetting interpreterSetting = interpreterSettingManager.getByName("test"); - interpreterSetting.getOption().setPerUser(InterpreterOption.SHARED); - - Interpreter interpreter1 = interpreterSetting.getDefaultInterpreter("user1", note1Id); - RemoteInterpreter remoteInterpreter1 = (RemoteInterpreter) interpreter1; - InterpreterContext context1 = InterpreterContext.builder() - .setNoteId("noteId") - .setParagraphId("paragraphId") - .build(); - remoteInterpreter1.interpret("hello", context1); - - assertEquals(1, interpreterSettingManager.getRecoveryStorage().restore().size()); - - interpreterSetting.close(); - assertEquals(0, interpreterSettingManager.getRecoveryStorage().restore().size()); - } - - @Test - void testMultipleInterpreterProcess() throws InterpreterException, IOException { - InterpreterSetting interpreterSetting = interpreterSettingManager.getByName("test"); - interpreterSetting.getOption().setPerUser(InterpreterOption.ISOLATED); - - Interpreter interpreter1 = interpreterSetting.getDefaultInterpreter("user1", note1Id); - RemoteInterpreter remoteInterpreter1 = (RemoteInterpreter) interpreter1; - InterpreterContext context1 = InterpreterContext.builder() - .setNoteId("noteId") - .setParagraphId("paragraphId") - .build(); - remoteInterpreter1.interpret("hello", context1); - assertEquals(1, interpreterSettingManager.getRecoveryStorage().restore().size()); - - Interpreter interpreter2 = interpreterSetting.getDefaultInterpreter("user2", note2Id); - RemoteInterpreter remoteInterpreter2 = (RemoteInterpreter) interpreter2; - InterpreterContext context2 = InterpreterContext.builder() - .setNoteId("noteId") - .setParagraphId("paragraphId") - .build(); - remoteInterpreter2.interpret("hello", context2); - - assertEquals(2, interpreterSettingManager.getRecoveryStorage().restore().size()); - - interpreterSettingManager.restart(interpreterSetting.getId(), "user1", note1Id); - assertEquals(1, interpreterSettingManager.getRecoveryStorage().restore().size()); - - interpreterSetting.close(); - assertEquals(0, interpreterSettingManager.getRecoveryStorage().restore().size()); - } - -} diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/realm/TestGroupResolver.java b/zeppelin-server/src/test/java/org/apache/zeppelin/realm/TestGroupResolver.java new file mode 100644 index 00000000000..1f4d104347f --- /dev/null +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/realm/TestGroupResolver.java @@ -0,0 +1,29 @@ +/* + * 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.zeppelin.realm; + +import java.util.Collections; +import java.util.Set; + +/** Group resolver used by Shiro integration tests that do not install optional plugins. */ +public class TestGroupResolver implements GroupResolver { + + @Override + public Set resolve(String principal) { + return Collections.emptySet(); + } +} diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/realm/jwt/KnoxJwtRealmTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/realm/jwt/KnoxJwtRealmTest.java index d9860f91073..5165d132a0b 100644 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/realm/jwt/KnoxJwtRealmTest.java +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/realm/jwt/KnoxJwtRealmTest.java @@ -26,6 +26,7 @@ import org.junit.jupiter.api.Test; import java.util.Date; +import java.util.Set; import static org.junit.jupiter.api.Assertions.*; @@ -96,6 +97,13 @@ void testValidateExpiration_WithPastExpiration_ShouldReturnFalse() throws Except assertFalse(result, "JWT token with past expiration should be rejected"); } + @Test + void testMapGroupPrincipalsUsesConfiguredResolver() { + knoxJwtRealm.setGroupResolver(principal -> Set.of(principal + "-role")); + + assertEquals(Set.of("alice-role"), knoxJwtRealm.mapGroupPrincipals("alice")); + } + // Note: Full token validation tests are omitted as they require complex setup // including certificate files and signature validation. The core expiration // validation logic is tested above through direct method calls. diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/recovery/RecoveryTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/recovery/RecoveryTest.java index 5a4f339feb0..f45fe2ad1aa 100644 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/recovery/RecoveryTest.java +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/recovery/RecoveryTest.java @@ -27,7 +27,7 @@ import org.apache.zeppelin.interpreter.InterpreterSetting; import org.apache.zeppelin.interpreter.InterpreterSettingManager; import org.apache.zeppelin.interpreter.ManagedInterpreterGroup; -import org.apache.zeppelin.interpreter.recovery.FileSystemRecoveryStorage; +import org.apache.zeppelin.interpreter.recovery.LocalRecoveryStorage; import org.apache.zeppelin.interpreter.recovery.StopInterpreter; import org.apache.zeppelin.notebook.Notebook; import org.apache.zeppelin.notebook.Paragraph; @@ -71,7 +71,7 @@ static void init() throws Exception { zepServer.copyBinDir(); zepServer.getZeppelinConfiguration().setProperty( ZeppelinConfiguration.ConfVars.ZEPPELIN_RECOVERY_STORAGE_CLASS.getVarName(), - FileSystemRecoveryStorage.class.getName()); + LocalRecoveryStorage.class.getName()); recoveryDir = Files.createTempDirectory("recovery").toFile(); zepServer.getZeppelinConfiguration().setProperty( ZeppelinConfiguration.ConfVars.ZEPPELIN_RECOVERY_DIR.getVarName(), diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/rest/AbstractTestRestApi.java b/zeppelin-server/src/test/java/org/apache/zeppelin/rest/AbstractTestRestApi.java index 6e3b4d8615a..d664ab7c0b5 100644 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/rest/AbstractTestRestApi.java +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/rest/AbstractTestRestApi.java @@ -91,6 +91,7 @@ public abstract class AbstractTestRestApi { "knoxJwtRealm.redirectParam = originalUrl\n" + "knoxJwtRealm.cookieName = hadoop-jwt\n" + "knoxJwtRealm.publicKeyPath = knox-sso.pem\n" + + "knoxJwtRealm.groupResolverClass = org.apache.zeppelin.realm.TestGroupResolver\n" + "authc = org.apache.zeppelin.realm.jwt.KnoxAuthenticationFilter\n" + "sessionManager = org.apache.shiro.web.session.mgt.DefaultWebSessionManager\n" + "securityManager.sessionManager = $sessionManager\n" + diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/server/ZeppelinServerPluginClassLoadingTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/server/ZeppelinServerPluginClassLoadingTest.java new file mode 100644 index 00000000000..92e93c7c798 --- /dev/null +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/server/ZeppelinServerPluginClassLoadingTest.java @@ -0,0 +1,126 @@ +/* + * 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.zeppelin.server; + +import static java.nio.charset.StandardCharsets.UTF_8; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertSame; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.io.IOException; +import java.io.InputStream; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.LinkedHashSet; +import java.util.Map; +import java.util.jar.JarEntry; +import java.util.jar.JarOutputStream; + +import org.eclipse.jetty.webapp.WebAppClassLoader; +import org.eclipse.jetty.webapp.WebAppContext; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; +import org.slf4j.Logger; + +class ZeppelinServerPluginClassLoadingTest { + @TempDir + Path tempDir; + + @Test + void loadsAllPluginJarsWhileKeepingSlf4jOnTheParentClassLoader() throws Exception { + Path firstJar = createJar( + "first.jar", Map.of("first-plugin-resource.txt", "first".getBytes(UTF_8))); + Path secondJar = createJar( + "second.jar", + Map.of( + "second-plugin-resource.txt", "second".getBytes(UTF_8), + "org/slf4j/Logger.class", readClassBytes(Logger.class))); + WebAppContext webapp = new WebAppContext(); + LinkedHashSet pluginClasspath = new LinkedHashSet<>(); + pluginClasspath.add(firstJar.toFile()); + pluginClasspath.add(secondJar.toFile()); + + ZeppelinServer.configurePluginClassLoading(webapp, pluginClasspath); + + ClassLoader parent = Thread.currentThread().getContextClassLoader(); + try (WebAppClassLoader classLoader = new WebAppClassLoader(parent, webapp)) { + assertEquals(2, classLoader.getURLs().length); + assertEquals("first", readResource(classLoader, "first-plugin-resource.txt")); + assertEquals("second", readResource(classLoader, "second-plugin-resource.txt")); + assertNotNull(classLoader.findResource("org/slf4j/Logger.class")); + assertSame(Logger.class, classLoader.loadClass(Logger.class.getName())); + } + } + + @Test + void onlyTreatsShiroObjectDeclarationsAsClassAssignments() { + Map assignments = ZeppelinServer.findShiroClassAssignments( + "[users]\n" + + "alice = secret.example\n" + + "[main]\n" + + "krbRealm = org.apache.zeppelin.realm.kerberos.KerberosRealm\n" + + "krbRealm.cookieDomain = domain.com\n" + + "knox.groupResolverClass = example.CustomGroupResolver\n" + + "[roles]\n" + + "analyst = example.Role\n"); + + assertEquals(1, assignments.size()); + assertEquals( + "org.apache.zeppelin.realm.kerberos.KerberosRealm", assignments.get("krbRealm")); + assertTrue(assignments.values().stream().noneMatch("domain.com"::equals)); + assertEquals( + "example.CustomGroupResolver", + ZeppelinServer.findNonBlankAssignment( + "[users]\n" + + "knox.groupResolverClass = ignored.UserValue\n" + + "[main]\n" + + "knox.groupResolverClass = example.CustomGroupResolver\n" + + "[roles]\n" + + "knox.groupResolverClass = ignored.RoleValue\n", + "knox.groupResolverClass")); + } + + private Path createJar(String fileName, Map entries) throws IOException { + Path jar = tempDir.resolve(fileName); + try (JarOutputStream output = new JarOutputStream(Files.newOutputStream(jar))) { + for (Map.Entry entry : entries.entrySet()) { + output.putNextEntry(new JarEntry(entry.getKey())); + output.write(entry.getValue()); + output.closeEntry(); + } + } + return jar; + } + + private static byte[] readClassBytes(Class clazz) throws IOException { + String resource = "/" + clazz.getName().replace('.', '/') + ".class"; + try (InputStream input = clazz.getResourceAsStream(resource)) { + if (input == null) { + throw new IOException("Class resource is missing: " + resource); + } + return input.readAllBytes(); + } + } + + private static String readResource(ClassLoader classLoader, String resource) throws IOException { + try (InputStream input = classLoader.getResourceAsStream(resource)) { + assertNotNull(input, "Resource is missing: " + resource); + return new String(input.readAllBytes(), UTF_8); + } + } +} From 3b89caea4743ef5758821ebb33bbb7ea3f836806 Mon Sep 17 00:00:00 2001 From: Jongyoul Lee Date: Tue, 4 Aug 2026 18:28:46 +0900 Subject: [PATCH 2/2] Fix Knox role test without Hadoop plugin --- .../apache/zeppelin/service/ShiroAuthenticationServiceTest.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/service/ShiroAuthenticationServiceTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/service/ShiroAuthenticationServiceTest.java index f82539e715d..8523900f475 100644 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/service/ShiroAuthenticationServiceTest.java +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/service/ShiroAuthenticationServiceTest.java @@ -36,6 +36,7 @@ import org.apache.shiro.util.LifecycleUtils; import org.apache.shiro.util.ThreadContext; import org.apache.zeppelin.conf.ZeppelinConfiguration; +import org.apache.zeppelin.realm.TestGroupResolver; import org.apache.zeppelin.realm.jwt.KnoxJwtRealm; import org.apache.zeppelin.service.shiro.AbstractShiroTest; import org.h2.jdbcx.JdbcDataSource; @@ -107,6 +108,7 @@ void testKnoxGetRoles() { setupPrincipalName("test"); KnoxJwtRealm realm = spy(new KnoxJwtRealm()); + realm.setGroupResolverClass(TestGroupResolver.class.getName()); LifecycleUtils.init(realm); Set testRoles = new HashSet(); testRoles.add("role1");