diff --git a/docs/deployment/migration-guide.md b/docs/deployment/migration-guide.md index d00b1a63a59..c58e01a9414 100644 --- a/docs/deployment/migration-guide.md +++ b/docs/deployment/migration-guide.md @@ -23,6 +23,7 @@ * Since Kyuubi 1.13, the support of Flink engine for Flink 1.17, 1.18 and 1.19 is removed. * Since Kyuubi 1.13, the support of Flink engine for Flink 2.0 is deprecated, and will be removed in the future. * Since Kyuubi 1.13, `kyuubi.server.redaction.regex` defaults to `(?i)secret|password|token|access[.]key` instead of being unset, so `kyuubi.server.conf.retrieveMode=REDACTED` (the default) redacts matching session-config keys/values out of the box; it also affects the command-line arguments Kyuubi logs for spawned engine processes. Set it to a different pattern to override. +* Since Kyuubi 1.13, the Spark authorization plugin is migrated from the Ranger plugin API to the Ranger 2.9 authorization API, and Apache Ranger 2.9.0 or above is required for both authorizer modes. See [Installing and Configuring Kyuubi Spark AuthZ Plugin](../security/authorization/spark/install.md) for how to set up the new API. ## Upgrading from Kyuubi 1.11 to 1.12 diff --git a/docs/security/authorization/spark/build.md b/docs/security/authorization/spark/build.md index abfebdb471b..5c02b35ab87 100644 --- a/docs/security/authorization/spark/build.md +++ b/docs/security/authorization/spark/build.md @@ -74,24 +74,16 @@ The available `spark.version`s are shown in the following table. The maven option `ranger.version` is used for specifying Ranger version to compile with and generate corresponding transitive dependencies. By default, it is always built with the latest `ranger.version` defined in kyuubi project main pom file. -Sometimes, it may be incompatible with other Ranger Admins, then you may need to build the plugin on your own targeting the Ranger Admin version you connect with. ```shell -build/mvn clean package -pl :kyuubi-spark-authz_2.12 -am -DskipTests -Dranger.version=2.4.0 +build/mvn clean package -pl :kyuubi-spark-authz_2.12 -am -DskipTests -Dranger.version=2.9.0 ``` -The available `ranger.version`s are shown in the following table. +The plugin is built on the Ranger 2.9 authorization API (`ranger-authz-api` and +`authz-remote`), which was introduced in Ranger 2.9.0, so only Ranger 2.9.0 and above +are supported. -| Ranger Version | Supported | Remark | -|:--------------:|:---------:|:------:| -| 2.6.x | √ | - | -| 2.5.x | √ | - | -| 2.4.x | √ | - | -| 2.3.x | √ | - | -| 2.2.x | √ | - | -| 2.1.x | √ | - | - -Currently, all ranger releases are supported. +Please use branch-1.12 or prior to build against Ranger versions prior to 2.9.0. ## Test with ScalaTest Maven plugin diff --git a/docs/security/authorization/spark/install.md b/docs/security/authorization/spark/install.md index 94419ff91a3..804bf267129 100644 --- a/docs/security/authorization/spark/install.md +++ b/docs/security/authorization/spark/install.md @@ -21,12 +21,12 @@ - [Apache Ranger](https://ranger.apache.org/) - This plugin works as a ranger rest client with Apache Ranger Admin server to do privilege check. - Thus, a ranger server need to be installed ahead and available to use. + The plugin talks to a Ranger PDP server (default) or Ranger Admin server to do privilege check. + A Ranger 2.9.0 server or above needs to be installed ahead and available to use. - Building(optional) - If your Ranger Admin or Spark distribution is not compatible with the official pre-built [artifact](https://mvnrepository.com/artifact/org.apache.kyuubi/kyuubi-spark-authz) in maven central. + If your Ranger or Spark distribution is not compatible with the official pre-built [artifact](https://mvnrepository.com/artifact/org.apache.kyuubi/kyuubi-spark-authz) in maven central. You need to [build](build.md) the plugin targeting the spark/ranger you are using by yourself. ## Install @@ -37,8 +37,75 @@ Use either the shaded jar `kyuubi-spark-authz-shaded_*.jar` or the `kyuubi-spark ## Configure +### Authorizer Modes + +The plugin supports two authorizer modes, selected by `ranger.authorizer.impl.class`: + +- `org.apache.ranger.authz.remote.RangerRemoteAuthorizer` (default) — Ranger PDP mode. + Authorization requests are sent to the Ranger PDP server via REST APIs. The plugin is a thin + client: policies are not downloaded to the client side, and access audits are recorded by + the PDP server. +- `org.apache.ranger.authz.embedded.RangerEmbeddedAuthorizer` — embedded mode, working as the + previous Ranger plugin did: policies are pulled from the Ranger admin server and access + requests are evaluated locally on the Spark driver side. + +:::{warning} +The security-critical configurations — the authorizer implementation +(`ranger.authorizer.impl.class`), the Ranger PDP server address +(`ranger.authz.remote.pdp.url`), the Ranger service name +(`ranger.plugin.spark.service.name`) and the Ranger PDP client authentication and +encryption settings (`ranger.authz.remote.authn.*`, `ranger.authz.remote.ssl.*`, +`ranger.authz.remote.header.*`) — are read from `ranger-spark-security.xml` only: +JVM system properties can neither set nor override them, and the plugin fails to +initialize when the Ranger PDP server address is missing for the Ranger PDP mode. +Other configurations with the `ranger.` or `xasecure.` prefix can still be overridden +by JVM system properties, so make sure tenant users cannot control the JVM options of +the engine process (e.g. `spark.driver.extraJavaOptions`). +::: + +#### Ranger PDP mode (default) + +- Create `ranger-spark-security.xml` in `$SPARK_HOME/conf` and add the following configurations + for pointing to the right Ranger PDP server. + +```xml + + + ranger.plugin.spark.service.name + a ranger service name, e.g. a ranger hive service name + + + + ranger.authz.remote.pdp.url + ranger pdp server address like https://ranger-pdp.org:8585 + + + +``` + +The PDP client supports authentication and encryption settings, configured with +`ranger.authz.remote.authn.type` (`header`, `jwt` or `kerberos`), +`ranger.authz.remote.authn.*`, `ranger.authz.remote.ssl.*` and +`ranger.authz.remote.header.*` properties. Refer to the +[Ranger client libraries](https://cwiki.apache.org/confluence/display/RANGER/Ranger+Client+Libraries) +for the full list. + +#### Embedded mode + +Set the authorizer implementation to the embedded one in `ranger-spark-security.xml`: + +```xml + + ranger.authorizer.impl.class + org.apache.ranger.authz.embedded.RangerEmbeddedAuthorizer + +``` + ### Settings for Connecting Ranger Admin +The settings below apply to the embedded mode, where the plugin pulls policies from +the Ranger admin server and evaluates access requests locally. + #### ranger-spark-security.xml - Create `ranger-spark-security.xml` in `$SPARK_HOME/conf` and add the following configurations @@ -76,7 +143,7 @@ Use either the shaded jar `kyuubi-spark-authz-shaded_*.jar` or the `kyuubi-spark ##### Using Macros in Row Level Filters -Macros are now supported for using user/group/tag in row filter expressions, introduced in [Ranger 2.3](https://cwiki.apache.org/confluence/display/RANGER/Apache+Ranger+2.3.0+-+Release+Notes). This feature helps significantly simplify row filter expressions by using user/group/tag's attributes instead of explicit conditions. Considering a user with an attribute `born_city` of value `Guangzhou `, the row filter condition as `city='${{USER.born_city}}'` will be transformed to `city='Guangzhou'` in execution plan. More supported macros and usage refer to [RANGER-3605](https://issues.apache.org/jira/browse/RANGER-3605) and [RANGER-3550](https://issues.apache.org/jira/browse/RANGER-3550). Add the following configs to `ranger-spark-security.xml` to enable UserStore Enricher required by macros. +Macros are supported for using user/group/tag in row filter expressions (embedded mode only), introduced in [Ranger 2.3](https://cwiki.apache.org/confluence/display/RANGER/Apache+Ranger+2.3.0+-+Release+Notes). This feature helps significantly simplify row filter expressions by using user/group/tag's attributes instead of explicit conditions. Considering a user with an attribute `born_city` of value `Guangzhou `, the row filter condition as `city='${{USER.born_city}}'` will be transformed to `city='Guangzhou'` in execution plan. More supported macros and usage refer to [RANGER-3605](https://issues.apache.org/jira/browse/RANGER-3605) and [RANGER-3550](https://issues.apache.org/jira/browse/RANGER-3550). Add the following configs to `ranger-spark-security.xml` to enable UserStore Enricher required by macros. ```xml @@ -94,7 +161,7 @@ Macros are now supported for using user/group/tag in row filter expressions, int ##### Showing all disallowed privileges -By default, Authz plugin checks required privileges one by one and throw the first unsatisfied privilege in exception. By setting `ranger.plugin.spark.authorize.in.single.call` to `true`, Authz plugin executes access checks in single call and throws all disallowed privileges in exception message. +By default, Authz plugin checks required privileges one by one and throw the first unsatisfied privilege in exception. By setting `ranger.plugin.spark.authorize.in.single.call` to `true`, Authz plugin executes access checks in single call and throws all disallowed privileges in exception message. This setting also reduces the number of authorization requests in Ranger PDP mode. ```xml @@ -106,8 +173,9 @@ By default, Authz plugin checks required privileges one by one and throw the fir #### ranger-spark-audit.xml -Create `ranger-spark-audit.xml` in `$SPARK_HOME/conf` and add the following configurations -to enable/disable auditing. +In the embedded mode, create `ranger-spark-audit.xml` in `$SPARK_HOME/conf` and add the following +configurations to enable/disable auditing. In Ranger PDP mode, access audits are recorded by +the PDP server, and this file is not required. ```xml diff --git a/extensions/spark/kyuubi-spark-authz-shaded/pom.xml b/extensions/spark/kyuubi-spark-authz-shaded/pom.xml index 42a5ee2ea4b..79a5adf41a1 100644 --- a/extensions/spark/kyuubi-spark-authz-shaded/pom.xml +++ b/extensions/spark/kyuubi-spark-authz-shaded/pom.xml @@ -48,15 +48,13 @@ org.apache.kyuubi:* - org.apache.ranger:* - - org.codehaus.jackson:* + org.apache.ranger:ranger-authz-api + org.apache.ranger:authz-remote com.fasterxml.jackson.core:* - com.fasterxml.jackson.module:* - com.fasterxml.jackson.jaxrs:* - com.sun.jersey:* - javax.ws.rs:jsr311-api - commons-collections:commons-collections + com.fasterxml.jackson.module:jackson-module-scala_* + org.apache.httpcomponents:httpclient + org.apache.httpcomponents:httpcore + org.apache.commons:commons-lang3 @@ -82,33 +80,17 @@ - - org.codehaus.jackson - ${kyuubi.shade.packageName}.org.codehaus.jackson - com.fasterxml.jackson ${kyuubi.shade.packageName}.com.fasterxml.jackson - com.sun.jersey - ${kyuubi.shade.packageName}.com.sun.jersey - - - com.sun.ws.rs.ext - ${kyuubi.shade.packageName}.com.sun.ws.rs.ext - - - javax.ws.rs - ${kyuubi.shade.packageName}.javax.ws.rs - - - com.kstruct.gethostname4j - ${kyuubi.shade.packageName}.com.kstruct.gethostname4j + org.apache.http + ${kyuubi.shade.packageName}.org.apache.http - org.apache.commons.collections - ${kyuubi.shade.packageName}.org.apache.commons.collections + org.apache.commons.lang3 + ${kyuubi.shade.packageName}.org.apache.commons.lang3 diff --git a/extensions/spark/kyuubi-spark-authz-shaded/src/main/resources/META-INF/LICENSE b/extensions/spark/kyuubi-spark-authz-shaded/src/main/resources/META-INF/LICENSE index 1e6d25e885e..c63740d7549 100644 --- a/extensions/spark/kyuubi-spark-authz-shaded/src/main/resources/META-INF/LICENSE +++ b/extensions/spark/kyuubi-spark-authz-shaded/src/main/resources/META-INF/LICENSE @@ -207,19 +207,12 @@ This project bundles some components that are licensed under the Apache License Version 2.0 -------------------------- -org.apache.ranger:ranger-plugins-common -org.apache.ranger:ranger-plugins-audit -org.codehaus.jackson:jackson-jaxrs -org.codehaus.jackson:jackson-core-asl -org.codehaus.jackson:jackson-mapper-asl -net.java.dev.jna:jna -net.java.dev.jna:jna-platform - -Common Development and Distribution License (CDDL) 1.1 ------------------------------------------------------- -com.sun.jersey:jersey-client -com.sun.jersey:jersey-core - -MIT license ------------ -com.kstruct:gethostname4j +com.fasterxml.jackson.core:jackson-annotations +com.fasterxml.jackson.core:jackson-core +com.fasterxml.jackson.core:jackson-databind +com.fasterxml.jackson.module:jackson-module-scala_* +org.apache.commons:commons-lang3 +org.apache.httpcomponents:httpclient +org.apache.httpcomponents:httpcore +org.apache.ranger:authz-remote +org.apache.ranger:ranger-authz-api diff --git a/extensions/spark/kyuubi-spark-authz/README.md b/extensions/spark/kyuubi-spark-authz/README.md index 9755db9d81c..1980762f8be 100644 --- a/extensions/spark/kyuubi-spark-authz/README.md +++ b/extensions/spark/kyuubi-spark-authz/README.md @@ -23,10 +23,15 @@ - [x] Row-level fine-grained authorization, a.k.a. Row-level filtering - [x] Data masking +The plugin supports two authorizer modes: Ranger PDP mode (default) which sends +authorization requests to a Ranger PDP server via REST APIs with a thin client, and +the embedded mode which pulls policies from the Ranger admin server and evaluates +requests locally. See `docs/security/authorization/spark/install.md` for details. + ## Build ```shell -build/mvn clean package -DskipTests -pl :kyuubi-spark-authz_2.12 -am -Dspark.version=3.5.6 -Dranger.version=2.6.0 +build/mvn clean package -DskipTests -pl :kyuubi-spark-authz_2.12 -am -Dspark.version=3.5.6 -Dranger.version=2.9.0 ``` ### Supported Apache Spark Versions @@ -44,11 +49,5 @@ build/mvn clean package -DskipTests -pl :kyuubi-spark-authz_2.12 -am -Dspark.ver `-Dranger.version=` -- [ ] 2.7.x -- [x] 2.6.x (default) -- [x] 2.5.x -- [x] 2.4.x -- [x] 2.3.x -- [x] 2.2.x -- [x] 2.1.x -- [ ] 2.0.x +The plugin is built on the Ranger 2.9 authorization API, so Ranger 2.9.0 and above are +the only supported versions. diff --git a/extensions/spark/kyuubi-spark-authz/pom.xml b/extensions/spark/kyuubi-spark-authz/pom.xml index 97e9ee01fb3..0602e1033fe 100644 --- a/extensions/spark/kyuubi-spark-authz/pom.xml +++ b/extensions/spark/kyuubi-spark-authz/pom.xml @@ -32,8 +32,7 @@ https://kyuubi.apache.org/ - 2.6.0 - 1.19.4 + 2.9.0 @@ -61,191 +60,37 @@ kyuubi-util-scala_${scala.binary.version} ${project.version} - - org.apache.ranger - ranger-plugins-common - ${ranger.version} - - - org.apache.ranger - ranger-plugin-classloader - - - org.apache.ranger - ranger-plugins-audit - - - - com.sun.jersey - jersey-bundle - - - log4j - log4j - - - ch.qos.logback - logback-classic - - - org.apache.commons - commons-configuration2 - - - commons-logging - commons-logging - - - org.apache.hadoop - hadoop-common - - - javax.ws.rs - jsr311-api - - - com.kstruct - gethostname4j - - - net.java.dev.jna - jna - - - net.java.dev.jna - jna-platform - - - - - - com.sun.jersey - jersey-client - ${jersey.client.version} - + org.apache.ranger - ranger-plugin-classloader + ranger-authz-api ${ranger.version} org.apache.ranger - ranger-plugins-audit + authz-remote ${ranger.version} - - - org.apache.ranger - ranger-plugins-cred - - - org.apache.kafka - * - - - org.apache.solr - solr-solrj - - - org.elasticsearch - * - - - org.elasticsearch.client - * - - - org.elasticsearch.plugin - * - - - org.apache.lucene - * - - - log4j - log4j - - - commons-lang - commons-lang - - - commons-logging - commons-logging - - - com.carrotsearch - hppc - - - org.apache.httpcomponents - * - - - org.apache.hive - hive-storage-api - - - org.apache.orc - orc-core - - - org.apache.hadoop - hadoop-common - - - com.google.guava - guava - - - joda-time - joda-time - - - org.apache.logging.log4j - * - - - - com.amazonaws - aws-java-sdk-bundle - - - com.amazonaws - aws-java-sdk-logs - - + org.apache.ranger - ranger-plugins-cred + authz-embedded ${ranger.version} + test + - org.apache.commons - commons-configuration2 - - - org.apache.hadoop - hadoop-common - - - log4j - log4j + org.slf4j + slf4j-reload4j - - org.eclipse.persistence - javax.persistence - - - org.eclipse.persistence - eclipselink + ch.qos.reload4j + reload4j @@ -280,11 +125,6 @@ provided - - commons-collections - commons-collections - - com.fasterxml.jackson.module jackson-module-scala_${scala.binary.version} diff --git a/extensions/spark/kyuubi-spark-authz/src/main/java/com/kstruct/gethostname4j/Hostname.java b/extensions/spark/kyuubi-spark-authz/src/main/java/com/kstruct/gethostname4j/Hostname.java deleted file mode 100644 index ac384326ac6..00000000000 --- a/extensions/spark/kyuubi-spark-authz/src/main/java/com/kstruct/gethostname4j/Hostname.java +++ /dev/null @@ -1,52 +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 com.kstruct.gethostname4j; - -import java.net.InetAddress; -import java.net.UnknownHostException; -import org.apache.commons.lang3.StringUtils; -import org.apache.commons.lang3.SystemUtils; - -// Alternative for RANGER-4125 to cut out JNA dependencies -public class Hostname { - - // The highest priority environment variable which allows user to - // set hostname for Ranger client - public static final String RANGER_CLIENT_HOSTNAME = "RANGER_CLIENT_HOSTNAME"; - - /** @return the hostname the of the current machine */ - public static String getHostname() { - String hostname = System.getenv(RANGER_CLIENT_HOSTNAME); - if (isValid(hostname)) return hostname; - - // Gets the host name from an environment variable - // (COMPUTERNAME on Windows, HOSTNAME elsewhere) - hostname = SystemUtils.getHostName(); - if (isValid(hostname)) return hostname; - - try { - return InetAddress.getLocalHost().getHostName(); - } catch (UnknownHostException rethrow) { - throw new RuntimeException(rethrow); - } - } - - private static boolean isValid(String hostname) { - return StringUtils.isNotBlank(hostname) && !"localhost".equals(hostname); - } -} diff --git a/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/AccessRequest.scala b/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/AccessRequest.scala index 8fc8028e683..2b862999317 100644 --- a/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/AccessRequest.scala +++ b/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/AccessRequest.scala @@ -17,81 +17,34 @@ package org.apache.kyuubi.plugin.spark.authz.ranger -import java.util.{HashMap => JHashMap, Set => JSet} -import java.util.Date - -import scala.collection.JavaConverters._ - import org.apache.hadoop.security.UserGroupInformation -import org.apache.ranger.plugin.policyengine.{RangerAccessRequestImpl, RangerPolicyEngine} import org.apache.kyuubi.plugin.spark.authz.OperationType.OperationType -import org.apache.kyuubi.plugin.spark.authz.ranger.AccessType._ -import org.apache.kyuubi.util.reflect.ReflectUtils._ +import org.apache.kyuubi.plugin.spark.authz.ranger.AccessType.AccessType -case class AccessRequest private (accessType: AccessType) extends RangerAccessRequestImpl +/** + * A request to authorize a user for accessing a resource with an access type. + * + * @param resource the resource to authorize + * @param user the name of the user to authorize + * @param userGroups the groups of the user + * @param opType the Spark SQL operation type requesting the access + * @param accessType the access type to authorize + */ +case class AccessRequest private[ranger] ( + resource: AccessResource, + user: String, + userGroups: Set[String], + opType: OperationType, + accessType: AccessType) object AccessRequest { + def apply( resource: AccessResource, user: UserGroupInformation, opType: OperationType, accessType: AccessType): AccessRequest = { - val userName = user.getShortUserName - val userGroups = getUserGroups(user) - val req = new AccessRequest(accessType) - req.setResource(resource) - req.setUser(userName) - req.setUserGroups(userGroups) - req.setAction(opType.toString) - try { - val roles = invokeAs[JSet[String]]( - SparkRangerAdminPlugin, - "getRolesFromUserAndGroups", - (classOf[String], userName), - (classOf[JSet[String]], userGroups)) - invokeAs[Unit](req, "setUserRoles", (classOf[JSet[String]], roles)) - } catch { - case _: Exception => - } - req.setAccessTime(new Date()) - accessType match { - case USE => req.setAccessType(RangerPolicyEngine.ANY_ACCESS) - case _ => req.setAccessType(accessType.toString.toLowerCase) - } - try { - val clusterName = invokeAs[String](SparkRangerAdminPlugin, "getClusterName") - invokeAs[Unit](req, "setClusterName", (classOf[String], clusterName)) - } catch { - case _: Exception => - } - req - } - - private def getUserGroupsFromUgi(user: UserGroupInformation): JSet[String] = { - user.getGroupNames.toSet.asJava - } - - private def getUserGroupsFromUserStore(user: UserGroupInformation): Option[JSet[String]] = { - try { - val storeEnricher = invokeAs[AnyRef](SparkRangerAdminPlugin, "getUserStoreEnricher") - val userStore = invokeAs[AnyRef](storeEnricher, "getRangerUserStore") - val userGroupMapping = - invokeAs[JHashMap[String, JSet[String]]](userStore, "getUserGroupMapping") - Some(userGroupMapping.get(user.getShortUserName)) - } catch { - case _: NoSuchMethodException => - None - } + AccessRequest(resource, user.getShortUserName, user.getGroupNames.toSet, opType, accessType) } - - private def getUserGroups(user: UserGroupInformation): JSet[String] = { - if (SparkRangerAdminPlugin.useUserGroupsFromUserStoreEnabled) { - getUserGroupsFromUserStore(user) - .getOrElse(getUserGroupsFromUgi(user)) - } else { - getUserGroupsFromUgi(user) - } - } - } diff --git a/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/AccessResource.scala b/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/AccessResource.scala index 70943187229..deab69be8f6 100644 --- a/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/AccessResource.scala +++ b/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/AccessResource.scala @@ -18,27 +18,167 @@ package org.apache.kyuubi.plugin.spark.authz.ranger import java.io.File -import java.util -import scala.language.implicitConversions +import scala.collection.JavaConverters._ -import org.apache.ranger.plugin.policyengine.RangerAccessResourceImpl +import org.apache.commons.lang3.StringUtils +import org.apache.ranger.authz.model.RangerResourceInfo -import org.apache.kyuubi.plugin.spark.authz.{ObjectType, PrivilegeObject} +import org.apache.kyuubi.plugin.spark.authz.{AccessControlException, ObjectType, PrivilegeObject} import org.apache.kyuubi.plugin.spark.authz.ObjectType._ import org.apache.kyuubi.plugin.spark.authz.OperationType.OperationType -class AccessResource private (val objectType: ObjectType, val catalog: Option[String]) - extends RangerAccessResourceImpl { - implicit def asString(obj: Object): String = if (obj != null) obj.asInstanceOf[String] else null - def getDatabase: String = getValue("database") - def getUdf: String = getValue("udf") - def getTable: String = getValue("table") - def getColumn: String = getValue("column") +/** + * A privilege object to authorize, which is converted to a Ranger resource + * (e.g. "table:default/src" or "column:default/src/id") in requests to + * the Ranger authorizer. + * + * @param objectType the type of the object + * @param database the database name, or null if not applicable + * @param table the table name, or null if not applicable + * @param column the column names joined by comma, or null if not applicable + * @param udf the function name, or null if not applicable + * @param uri the uri path, or null if not applicable + * @param owner the owner of the object, if any + * @param catalog the catalog name, if any + */ +case class AccessResource private[ranger] ( + objectType: ObjectType, + database: String, + table: String, + column: String, + udf: String, + uri: String, + owner: Option[String], + catalog: Option[String]) { + + def getDatabase: String = database + def getUdf: String = udf + def getTable: String = table + def getColumn: String = column + def getColumns: Seq[String] = { - val columnStr = getColumn - if (columnStr == null) Nil else columnStr.split(",").filter(_.nonEmpty) + if (column == null) Nil else column.split(",").filter(_.nonEmpty) } + + def getOwnerUser: String = owner.orNull + + /** + * The path-like representation of this resource used in error messages, + * e.g. "default/src" for a table, "default/src/id" for a column. + */ + def getAsString: String = objectType match { + case COLUMN => + Seq(database, table, column).filter(_ != null).mkString("/") + case FUNCTION => + Seq(database, udf).filter(_ != null).mkString("/") + case URI => + // the uri is matched against both the exact path and the path with a + // trailing slash, as the legacy plugin did + val path = Option(uri).map(_.stripSuffix(File.separator)).getOrElse("") + s"[$path, $path/]" + case _ => + Seq(database, table).filter(_ != null).mkString("/") + } + + private[ranger] def toResourceInfos: Seq[RangerResourceInfo] = { + val attributes = owner.map(o => java.util.Collections.singletonMap("OWNER", o: AnyRef)).orNull + objectType match { + case DATABASE => + Seq(new RangerResourceInfo( + s"database:${requireRrnComponent(database, "database")}", + null, + null, + attributes)) + case FUNCTION => + // An unqualified function reference (e.g. a built-in or temporary function) has no + // database. The legacy plugin left the database blank, and blank values in the legacy + // resource matched the wildcard values in policies, so keep the wildcard marker for + // the blank database. Blank names in other resource levels indicate a broken command + // extraction, which matched no policy in the legacy resource, so deny them. + val db = if (StringUtils.isBlank(database)) "*" else escapeRrnMetaChars(database) + Seq(new RangerResourceInfo( + s"udf:$db/${requireRrnComponent(udf, "udf")}", + null, + null, + attributes)) + case COLUMN => + val columns = getColumns + // all the column requests need the database and the table + val parent = + s"column:${requireRrnComponent(database, "database")}" + + s"/${requireRrnComponent(table, "table")}" + if (columns.length == 1) { + Seq(new RangerResourceInfo( + s"$parent/${requireRrnComponent(columns.head, "column")}", + null, + null, + attributes)) + } else if (columns.isEmpty) { + Seq(new RangerResourceInfo(parent, null, null, attributes)) + } else { + val subResources = columns + .map(col => s"column:${requireRrnComponent(col, "column")}") + .toSet.asJava + Seq(new RangerResourceInfo(parent, subResources, null, attributes)) + } + case URI => + // Url policies may be written with or without a trailing slash, and the legacy + // plugin matched a uri against both the exact path and the path with a trailing + // slash. The RRN request carries a single resource value set, so the two variants + // are returned and authorized as alternatives by separate requests. + val path = requireRrnComponent( + Option(uri).map(_.stripSuffix(File.separator)).orNull, + "uri") + Seq( + new RangerResourceInfo(s"url:$path", null, null, attributes), + new RangerResourceInfo(s"url:$path/", null, null, attributes)) + case VIEW => + // A local temporary view has no database in the SHOW TABLES output. The legacy + // plugin left the database blank, and blank values in the legacy resource matched + // the wildcard values in policies, so keep the wildcard marker for the blank + // database, as the unqualified function reference case does. + val db = if (StringUtils.isBlank(database)) "*" else escapeRrnMetaChars(database) + Seq(new RangerResourceInfo( + s"table:$db/${requireRrnComponent(table, "table")}", + null, + null, + attributes)) + case _ => + Seq(new RangerResourceInfo( + s"table:${requireRrnComponent(database, "database")}" + + s"/${requireRrnComponent(table, "table")}", + null, + null, + attributes)) + } + } + + /** + * Escapes the RRN metacharacters in a resource name, so that a name containing + * them is parsed as a single resource level instead of breaking the resource + * hierarchy, e.g. a table named "a/b" becomes "a\/b" in the resource name + * "table:default/a\/b". The escape format follows RangerResourceNameParser: + * a backslash escapes the next character, so "\\" stands for a literal + * backslash and "\/" stands for a literal separator. + */ + private def escapeRrnMetaChars(value: String): String = + value.replace("\\", "\\\\").replace("/", "\\/") + + /** + * Returns the RRN component for the given resource name, throwing an access + * control exception for a blank name. The RRN parser rejects blank resource + * values, and the legacy plugin had no value for a blank name, which matched + * no policy, so the access is denied rather than being evaluated against the + * wildcard marker. + */ + private def requireRrnComponent(value: String, name: String): String = + if (StringUtils.isBlank(value)) { + throw new AccessControlException( + s"Access denied: invalid [$objectType] resource, blank $name") + } else { + escapeRrnMetaChars(value) + } } object AccessResource { @@ -49,35 +189,41 @@ object AccessResource { secondLevelResource: String, thirdLevelResource: String, owner: Option[String] = None, - catalog: Option[String] = None): AccessResource = { - val resource = new AccessResource(objectType, catalog) - - resource.objectType match { - case DATABASE => resource.setValue("database", firstLevelResource) - case FUNCTION => - resource.setValue("database", Option(firstLevelResource).getOrElse("")) - resource.setValue("udf", secondLevelResource) - case COLUMN => - resource.setValue("database", firstLevelResource) - resource.setValue("table", secondLevelResource) - resource.setValue("column", thirdLevelResource) - case TABLE | VIEW | INDEX => - resource.setValue("database", firstLevelResource) - resource.setValue("table", secondLevelResource) - case URI => - val objectList = new util.ArrayList[String] - Option(firstLevelResource) - .filter(_.nonEmpty) - .foreach { path => - val s = path.stripSuffix(File.separator) - objectList.add(s) - objectList.add(s + File.separator) - } - resource.setValue("url", objectList) - } - resource.setServiceDef(SparkRangerAdminPlugin.getServiceDef) - owner.foreach(resource.setOwnerUser) - resource + catalog: Option[String] = None): AccessResource = objectType match { + case DATABASE => + new AccessResource(DATABASE, firstLevelResource, null, null, null, null, owner, catalog) + case FUNCTION => + new AccessResource( + FUNCTION, + Option(firstLevelResource).getOrElse(""), + null, + null, + secondLevelResource, + null, + owner, + catalog) + case COLUMN => + new AccessResource( + COLUMN, + firstLevelResource, + secondLevelResource, + thirdLevelResource, + null, + null, + owner, + catalog) + case TABLE | VIEW | INDEX => + new AccessResource( + objectType, + firstLevelResource, + secondLevelResource, + null, + null, + null, + owner, + catalog) + case URI => + new AccessResource(URI, null, null, null, null, firstLevelResource, owner, catalog) } def apply( diff --git a/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/RangerConfigProvider.scala b/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/RangerConfigProvider.scala deleted file mode 100644 index 05d8cc64f40..00000000000 --- a/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/RangerConfigProvider.scala +++ /dev/null @@ -1,46 +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.kyuubi.plugin.spark.authz.ranger - -import org.apache.hadoop.conf.Configuration - -import org.apache.kyuubi.plugin.spark.authz.util.AuthZUtils.isRanger21orGreater -import org.apache.kyuubi.util.reflect.ReflectUtils.invokeAs - -trait RangerConfigProvider { - - /** - * Get plugin config of different Ranger versions - * - * @return instance of - * org.apache.ranger.authorization.hadoop.config.RangerPluginConfig - * for Ranger 2.1 and above, - * or instance of - * org.apache.ranger.authorization.hadoop.config.RangerConfiguration - * for Ranger 2.0 and below - */ - val getRangerConf: Configuration = { - if (isRanger21orGreater) { - // for Ranger 2.1+ - invokeAs(this, "getConfig") - } else { - // for Ranger 2.0 and below - invokeAs("org.apache.ranger.authorization.hadoop.config.RangerConfiguration", "getInstance") - } - } -} diff --git a/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/RuleAuthorization.scala b/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/RuleAuthorization.scala index 676d12cf11e..c63b547e42d 100644 --- a/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/RuleAuthorization.scala +++ b/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/RuleAuthorization.scala @@ -19,10 +19,9 @@ package org.apache.kyuubi.plugin.spark.authz.ranger import scala.collection.mutable -import org.apache.ranger.plugin.policyengine.RangerAccessRequest -import org.apache.ranger.plugin.util.RangerPerfTracer import org.apache.spark.sql.SparkSession import org.apache.spark.sql.catalyst.plans.logical.LogicalPlan +import org.slf4j.LoggerFactory import org.apache.kyuubi.plugin.spark.authz._ import org.apache.kyuubi.plugin.spark.authz.ObjectType._ @@ -32,20 +31,14 @@ import org.apache.kyuubi.plugin.spark.authz.rule.Authorization import org.apache.kyuubi.plugin.spark.authz.util.AuthZUtils._ case class RuleAuthorization(spark: SparkSession) extends Authorization(spark) { - private val PERF_SPARKAUTH_REQUEST_LOG = - RangerPerfTracer.getPerfLogger("sparkauth.request") + // the perf logger used by the legacy Ranger plugin (RangerPerfTracer) to trace + // access checks; keep the same logger name for existing perf-log consumers + final private val PERF_LOGGER = + LoggerFactory.getLogger("org.apache.ranger.perf.sparkauth.request") override def checkPrivileges(spark: SparkSession, plan: LogicalPlan): Unit = { - val perf = if (RangerPerfTracer.isPerfTraceEnabled(PERF_SPARKAUTH_REQUEST_LOG)) { - RangerPerfTracer.getPerfTracer( - PERF_SPARKAUTH_REQUEST_LOG, - "RuleAuthorization.checkPrivileges()") - } else { - null - } - + val start = System.nanoTime try { - val auditHandler = new SparkRangerAuditHandler val ugi = getAuthzUgi(spark.sparkContext) val (inputs, outputs, opType) = PrivilegesBuilder.build(plan, spark) @@ -69,7 +62,7 @@ case class RuleAuthorization(spark: SparkSession) extends Authorization(spark) { addAccessRequest(outputs, isInput = false) val requestArrays = requests.map { request => - val resource = request.getResource.asInstanceOf[AccessResource] + val resource = request.resource resource.objectType match { case ObjectType.COLUMN if resource.getColumns.nonEmpty => resource.getColumns.map { col => @@ -81,21 +74,24 @@ case class RuleAuthorization(spark: SparkSession) extends Authorization(spark) { col, Option(resource.getOwnerUser), resource.catalog) - AccessRequest(cr, ugi, opType, request.accessType).asInstanceOf[RangerAccessRequest] + AccessRequest(cr, ugi, opType, request.accessType) } case _ => Seq(request) } }.toSeq if (authorizeInSingleCall) { - verify(requestArrays.flatten, auditHandler) + verify(requestArrays.flatten) } else { requestArrays.flatten.foreach { req => - verify(Seq(req), auditHandler) + verify(Seq(req)) } } } finally { - RangerPerfTracer.log(perf) + if (PERF_LOGGER.isDebugEnabled) { + val elapsed = System.nanoTime - start + PERF_LOGGER.debug(s"RuleAuthorization.checkPrivileges() took ${(elapsed / 1000000d)}ms") + } } } } diff --git a/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/RuleFunctionAuthorization.scala b/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/RuleFunctionAuthorization.scala index 0701bd26379..faa8d04f882 100644 --- a/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/RuleFunctionAuthorization.scala +++ b/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/RuleFunctionAuthorization.scala @@ -19,7 +19,6 @@ package org.apache.kyuubi.plugin.spark.authz.ranger import scala.collection.mutable -import org.apache.ranger.plugin.policyengine.RangerAccessRequest import org.apache.spark.sql.SparkSession import org.apache.spark.sql.catalyst.plans.logical.LogicalPlan @@ -37,7 +36,6 @@ case class RuleFunctionAuthorization(spark: SparkSession) extends (LogicalPlan = return } - val auditHandler = new SparkRangerAuditHandler val ugi = getAuthzUgi(spark.sparkContext) val (inputs, _, opType) = PrivilegesBuilder.buildFunctions(plan, spark) @@ -59,14 +57,11 @@ case class RuleFunctionAuthorization(spark: SparkSession) extends (LogicalPlan = addAccessRequest(inputs, isInput = true) - val requestSeq: Seq[RangerAccessRequest] = - requests.map(_.asInstanceOf[RangerAccessRequest]).toSeq - if (authorizeInSingleCall) { - verify(requestSeq, auditHandler) + verify(requests.toSeq) } else { - requestSeq.foreach { req => - verify(Seq(req), auditHandler) + requests.foreach { req => + verify(Seq(req)) } } } diff --git a/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/SparkRangerAdminPlugin.scala b/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/SparkRangerAdminPlugin.scala index 2744567a0fa..37a9a4f0673 100644 --- a/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/SparkRangerAdminPlugin.scala +++ b/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/SparkRangerAdminPlugin.scala @@ -17,20 +17,86 @@ package org.apache.kyuubi.plugin.spark.authz.ranger +import java.util.Properties + import scala.collection.JavaConverters._ import scala.collection.mutable.{ArrayBuffer, LinkedHashMap} +import org.apache.hadoop.conf.Configuration import org.apache.hadoop.util.ShutdownHookManager -import org.apache.ranger.plugin.policyengine.RangerAccessRequest -import org.apache.ranger.plugin.service.RangerBasePlugin +import org.apache.ranger.authz.api.{RangerAuthorizer, RangerAuthorizerFactory} +import org.apache.ranger.authz.model._ +import org.apache.ranger.authz.model.RangerAccessContext.CONTEXT_INFO_CLUSTER_NAME +import org.apache.ranger.authz.model.RangerAuthzResult.AccessDecision import org.slf4j.LoggerFactory import org.apache.kyuubi.plugin.spark.authz.AccessControlException +import org.apache.kyuubi.plugin.spark.authz.ObjectType +import org.apache.kyuubi.plugin.spark.authz.ranger.AccessType._ -object SparkRangerAdminPlugin extends RangerBasePlugin("spark", "sparkSql") - with RangerConfigProvider { +/** + * The entry point of Ranger authorization for Spark SQL, which delegates + * authorization requests to a Ranger authorizer: + *
    + *
  • org.apache.ranger.authz.remote.RangerRemoteAuthorizer (default): + * sends requests to a Ranger PDP server via REST APIs, with a thin client + * that does not download policies to the client side.
  • + *
  • org.apache.ranger.authz.embedded.RangerEmbeddedAuthorizer: + * evaluates requests locally against the policies pulled from + * the Ranger admin server, as the previous Ranger plugin did.
  • + *
+ * + * The authorizer implementation is selected by `ranger.authorizer.impl.class`. + */ +object SparkRangerAdminPlugin { final private val LOG = LoggerFactory.getLogger(getClass) + /** + * The service type of the Spark SQL service definition registered in + * Ranger admin. + */ + final val SERVICE_TYPE: String = "spark" + + final private val APP_TYPE: String = "sparkSql" + + final private val KEY_IMPL_CLASS = "ranger.authorizer.impl.class" + final private val KEY_PDP_URL = "ranger.authz.remote.pdp.url" + final private val KEY_SERVICE_NAME = s"ranger.plugin.$SERVICE_TYPE.service.name" + final private val KEY_CLUSTER_NAME = s"ranger.plugin.$SERVICE_TYPE.access.cluster.name" + final private val KEY_AUTHORIZE_IN_SINGLE_CALL = + s"ranger.plugin.$SERVICE_TYPE.authorize.in.single.call" + + final private val REMOTE_AUTHORIZER_CLASS: String = + classOf[org.apache.ranger.authz.remote.RangerRemoteAuthorizer].getName + + /** + * The security-critical configurations which are read from the Ranger configuration + * resources only: JVM system properties can neither set nor override them, so a tenant + * who can control the JVM options of the engine process cannot point the authorization + * to a server they control or bypass the policies. + */ + final private val ADMIN_ONLY_KEYS: Seq[String] = + Seq(KEY_IMPL_CLASS, KEY_PDP_URL, KEY_SERVICE_NAME) + + final private val ADMIN_ONLY_KEY_PREFIXES: Seq[String] = Seq( + "ranger.authz.remote.authn.", + "ranger.authz.remote.header.", + "ranger.authz.remote.ssl.") + + final private val CONFIG_PREFIXES = Seq("ranger.", "xasecure.") + + final private val CONFIG_RESOURCES = Seq( + s"ranger-$SERVICE_TYPE-security.xml", + s"ranger-$SERVICE_TYPE-audit.xml") + + private var authorizer: RangerAuthorizer = null + + /** + * The Ranger configurations loaded from the Hadoop configuration resources, + * e.g. ranger-spark-security.xml. + */ + private[ranger] var config: Configuration = null + /** * For a Spark SQL query, it may contain 0 or more privilege objects to verify, e.g. a typical * JOIN operator may have two tables and their columns to verify. @@ -38,93 +104,267 @@ object SparkRangerAdminPlugin extends RangerBasePlugin("spark", "sparkSql") * This configuration controls whether to verify the privilege objects in single call or * to verify them one by one. */ - def authorizeInSingleCall: Boolean = getRangerConf.getBoolean( - s"ranger.plugin.${getServiceType}.authorize.in.single.call", - false) + def authorizeInSingleCall: Boolean = config.getBoolean(KEY_AUTHORIZE_IN_SINGLE_CALL, false) - /** - * This configuration controls whether to override user's usergroups - * by the mapping fetched from Ranger's UserStore. - * - * It relies on Ranger's UserStore is a feature supported since Ranger 2.1. - * - * If true, user bound usergroups will be looked up in in Ranger's UserStore - * and the usergroups of AccessRequest is overriden. - * - * Please make sure configs in Ranger set properly: - * 1. set `ranger.plugin.spark.enable.implicit.userstore.enricher` to true - * 2. set cache path for UserStore in `ranger.plugin.hive.policy.cache.dir` - * 3. at least one condition of policies containing scripts, e.g. {{USER.attr}} in row-filter - */ - def useUserGroupsFromUserStoreEnabled: Boolean = getRangerConf.getBoolean( - s"ranger.plugin.$getServiceType.use.usergroups.from.userstore.enabled", - false) + def getServiceType: String = SERVICE_TYPE + + private def serviceName: String = config.get(KEY_SERVICE_NAME) /** - * plugin initialization + * authorizer initialization * with cleanup shutdown hook registered */ def initialize(): Unit = { - this.init() - registerCleanupShutdownHook(this) + ensureConfig() + val props = new Properties + config.iterator.asScala + .filter(entry => CONFIG_PREFIXES.exists(entry.getKey.startsWith)) + .foreach(entry => props.put(entry.getKey, entry.getValue)) + System.getProperties.asScala + .filter { case (key, _) => + CONFIG_PREFIXES.exists(key.startsWith) && !isAdminOnlyKey(key) + } + .foreach { case (key, value) => props.put(key, value) } + validateAdminOnlyKeys(props) + initialize(props) + } + + private def isAdminOnlyKey(key: String): Boolean = + ADMIN_ONLY_KEYS.contains(key) || ADMIN_ONLY_KEY_PREFIXES.exists(key.startsWith) + + /** + * Fails fast when the Ranger configuration resources lack a configuration the + * authorizer needs: the authorizer implementation defaults to the Ranger PDP mode, + * which requires the Ranger PDP server address. The configuration cannot fall back + * to JVM system properties, which tenant users may control. + */ + private def validateAdminOnlyKeys(props: Properties): Unit = { + val implClass = props.getProperty(KEY_IMPL_CLASS, REMOTE_AUTHORIZER_CLASS) + if (implClass == REMOTE_AUTHORIZER_CLASS && props.getProperty(KEY_PDP_URL) == null) { + throw new IllegalArgumentException( + s"$KEY_PDP_URL must be configured in ${CONFIG_RESOURCES.mkString(", ")} and " + + "cannot be set by JVM system properties") + } + } + + private def ensureConfig(): Unit = { + if (config == null) { + config = new Configuration + CONFIG_RESOURCES.foreach(addResourceIfReadable) + } + } + + private def addResourceIfReadable(resource: String): Unit = { + val loader = Thread.currentThread().getContextClassLoader + val url = if (loader != null) loader.getResource(resource) else null + if (url != null) { + config.addResource(url) + } else { + config.addResource(resource) + } + } + + private[ranger] def initialize(props: Properties): Unit = synchronized { + if (authorizer == null) { + props.put("ranger.authz.app.type", APP_TYPE) + authorizer = RangerAuthorizerFactory.createAuthorizer(props) + authorizer.init() + registerCleanupShutdownHook(authorizer) + LOG.info( + s"initialized ranger authorizer, service: $serviceName, " + + s"impl: ${authorizer.getClass.getName}") + } } /** - * register shutdown hook for plugin cleanup + * Reset the authorizer to the uninitialized state, so that the next + * [[initialize]] call re-creates it. Intended for tests only. */ - private def registerCleanupShutdownHook(plugin: RangerBasePlugin): Unit = { + private[ranger] def reset(): Unit = synchronized { + if (authorizer != null) { + try authorizer.close() + catch { + case e: Exception => LOG.warn("failed to close ranger authorizer", e) + } + authorizer = null + } + } + + private def registerCleanupShutdownHook(authorizer: RangerAuthorizer): Unit = { ShutdownHookManager.get().addShutdownHook( () => { - if (plugin != null) { - LOG.info(s"clean up ranger plugin, appId: ${plugin.getAppId}") - plugin.cleanup() - plugin.getAuditProviderFactory.shutdown() + if (authorizer != null) { + LOG.info(s"clean up ranger authorizer, impl: ${authorizer.getClass.getName}") + try authorizer.close() + catch { + case e: Exception => LOG.warn("failed to close ranger authorizer", e) + } } }, Integer.MAX_VALUE) } + private def checkInitialized(): RangerAuthorizer = { + if (authorizer == null) { + initialize() + } + authorizer + } + + private def userInfo(req: AccessRequest): RangerUserInfo = + new RangerUserInfo(req.user, null, req.userGroups.asJava, null) + + private def permissionOf(accessType: AccessType): String = accessType match { + // the Ranger any-access marker (RangerPolicyEngine.ANY_ACCESS) allows any access + // type matched by policies, e.g. for SHOW commands + case USE => "_any" + case _ => accessType.toString.toLowerCase + } + + private def context(): RangerAccessContext = { + val additionalInfo = new java.util.HashMap[String, AnyRef] + val clusterName = config.get(KEY_CLUSTER_NAME) + if (clusterName != null) { + additionalInfo.put(CONTEXT_INFO_CLUSTER_NAME, clusterName) + } + val ctx = new RangerAccessContext(SERVICE_TYPE, serviceName) + ctx.setAccessTime(System.currentTimeMillis) + ctx.setAdditionalInfo(additionalInfo) + ctx + } + + private def toAuthzRequest( + req: AccessRequest, + resourceInfo: RangerResourceInfo): RangerAuthzRequest = { + val access = new RangerAccessInfo( + resourceInfo, + req.opType.toString, + java.util.Collections.singleton(permissionOf(req.accessType))) + new RangerAuthzRequest(userInfo(req), access, context()) + } + + private def toAuthzRequest(req: AccessRequest): RangerAuthzRequest = + toAuthzRequest(req, req.resource.toResourceInfos.head) + + /** + * batch verifying RangerAccessRequests + * and throws exception with all disallowed privileges + * for accessType and resources + */ + def verify(requests: Seq[AccessRequest]): Unit = { + if (requests.nonEmpty) { + val authorizer = checkInitialized() + val user = userInfo(requests.head) + val ctx = context() + // a uri resource has two variants, the exact path and the path with a trailing + // slash, which are authorized as alternatives by separate requests; the single + // call authorizes a fixed set of accesses conjunctively, so it is used only when + // the requests contain no uri resource + val hasUri = requests.exists(_.resource.objectType == ObjectType.URI) + val results = if (authorizeInSingleCall && !hasUri) { + val accesses = requests.map(toAuthzRequest(_).getAccess).asJava + (for { + multiResult <- + Option(authorizer.authorize(new RangerMultiAuthzRequest(user, accesses, ctx))) + accessResults <- Option(multiResult.getAccesses) + } yield accessResults.asScala.map(Option(_)).toSeq).getOrElse(Seq.empty) + } else { + requests.map { req => + val variantResults = req.resource.toResourceInfos.map { resourceInfo => + Option(authorizer.authorize(toAuthzRequest(req, resourceInfo))) + } + variantResults.find(_.exists(_.getDecision == AccessDecision.ALLOW)) + .getOrElse(variantResults.head) + } + } + + // a missing result means the authorizer could not evaluate the request, + // which is denied instead of skipped + val indices = results.padTo(requests.length, None).zipWithIndex.collect { + case (result, idx) if result.forall(_.getDecision != AccessDecision.ALLOW) => idx + } + if (indices.nonEmpty) { + val accessTypeToResource = + indices.foldLeft(LinkedHashMap.empty[String, ArrayBuffer[String]])((m, idx) => { + val req = requests(idx) + val accessType = permissionOf(req.accessType) + val resource = req.resource.getAsString + m.getOrElseUpdate(accessType, ArrayBuffer.empty[String]) + .append(resource) + m + }) + val errorMsg = accessTypeToResource + .map { case (accessType, resources) => + s"[$accessType] ${resources.mkString("privilege on [", ",", "]")}" + }.mkString(", ") + throw new AccessControlException( + s"Permission denied: user [${requests.head.user}] does not have $errorMsg") + } + } + } + + private def requireResult(result: RangerAuthzResult, req: AccessRequest): RangerAuthzResult = + if (result == null) { + // a missing result means the authorizer could not evaluate the request, + // which must not be treated as allowed + throw new AccessControlException( + s"Permission denied: no authorization result for user [${req.user}], " + + s"resource [${req.resource.getAsString}]") + } else { + result + } + def getFilterExpr(req: AccessRequest): Option[String] = { - val result = evalRowFilterPolicies(req, null) + val result = requireResult(checkInitialized().authorize(toAuthzRequest(req)), req) Option(result) - .filter(_.isRowFilterEnabled) + .filter(_.getDecision == AccessDecision.ALLOW) + .flatMap { r => + Option(r.getPermissions.get(permissionOf(req.accessType))) + } + .map(_.getRowFilter) + .filter(rf => rf != null && rf.getFilterExpr != null && rf.getFilterExpr.nonEmpty) .map(_.getFilterExpr) - .filter(fe => fe != null && fe.nonEmpty) } def getMaskingExpr(req: AccessRequest): Option[String] = { - val col = req.getResource.asInstanceOf[AccessResource].getColumn - val result = evalDataMaskPolicies(req, null) - Option(result).filter(_.isMaskEnabled).map { res => - if ("MASK_NULL".equalsIgnoreCase(res.getMaskType)) { - "NULL" - } else if ("CUSTOM".equalsIgnoreCase(result.getMaskType)) { - val maskVal = res.getMaskedValue - if (maskVal == null) { + val col = req.resource.getColumn + val result = requireResult(checkInitialized().authorize(toAuthzRequest(req)), req) + Option(result) + .filter(_.getDecision == AccessDecision.ALLOW) + .flatMap { r => + Option(r.getPermissions.get(permissionOf(req.accessType))) + } + .map(_.getDataMask) + .filter(dm => dm != null && dm.getMaskType != null) + .map { dm => + val maskType = dm.getMaskType + if ("MASK_NULL".equalsIgnoreCase(maskType)) { "NULL" + } else if ("CUSTOM".equalsIgnoreCase(maskType)) { + val maskVal = dm.getMaskedValue + if (maskVal == null) { + "NULL" + } else { + s"${maskVal.replace("{col}", col)}" + } } else { - s"${maskVal.replace("{col}", col)}" + maskType match { + case "MASK" => regexp_replace(col) + case "MASK_SHOW_FIRST_4" => + regexp_replace(col, hasLen = true) + case "MASK_SHOW_LAST_4" => + val left = regexp_replace(s"left($col, length($col) - 4)") + s"concat($left, right($col, 4))" + case "MASK_HASH" => s"md5(cast($col as string))" + case "MASK_DATE_SHOW_YEAR" => s"date_trunc('YEAR', $col)" + case _ => Option(dm.getMaskedValue) + .filter(_.nonEmpty) + .map(maskedValue => s"${maskedValue.replace("{col}", col)}") + .orNull + } } - } else if (result.getMaskTypeDef != null) { - result.getMaskTypeDef.getName match { - case "MASK" => regexp_replace(col) - case "MASK_SHOW_FIRST_4" => - regexp_replace(col, hasLen = true) - case "MASK_SHOW_LAST_4" => - val left = regexp_replace(s"left($col, length($col) - 4)") - s"concat($left, right($col, 4))" - case "MASK_HASH" => s"md5(cast($col as string))" - case "MASK_DATE_SHOW_YEAR" => s"date_trunc('YEAR', $col)" - case _ => result.getMaskTypeDef.getTransformer match { - case transformer if transformer != null && transformer.nonEmpty => - s"${transformer.replace("{col}", col)}" - case _ => null - } - } - } else { - null } - } + .filter(_ != null) } private def regexp_replace(expr: String, hasLen: Boolean = false): String = { @@ -137,38 +377,11 @@ object SparkRangerAdminPlugin extends RangerBasePlugin("spark", "sparkSql") } /** - * batch verifying RangerAccessRequests - * and throws exception with all disallowed privileges - * for accessType and resources + * verifying whether the user has any access to the resource, + * used for filtering outputs of SHOW commands */ - def verify( - requests: Seq[RangerAccessRequest], - auditHandler: SparkRangerAuditHandler): Unit = { - if (requests.nonEmpty) { - val results = SparkRangerAdminPlugin.isAccessAllowed(requests.asJava, auditHandler) - if (results != null) { - val indices = results.asScala.zipWithIndex.filter { case (result, idx) => - result != null && !result.getIsAllowed - }.map(_._2) - if (indices.nonEmpty) { - val user = requests.head.getUser - val accessTypeToResource = - indices.foldLeft(LinkedHashMap.empty[String, ArrayBuffer[String]])((m, idx) => { - val req = requests(idx) - val accessType = req.getAccessType - val resource = req.getResource.getAsString - m.getOrElseUpdate(accessType, ArrayBuffer.empty[String]) - .append(resource) - m - }) - val errorMsg = accessTypeToResource - .map { case (accessType, resources) => - s"[$accessType] ${resources.mkString("privilege on [", ",", "]")}" - }.mkString(", ") - throw new AccessControlException( - s"Permission denied: user [$user] does not have $errorMsg") - } - } - } + def isAccessAllowed(req: AccessRequest): Boolean = { + val result = checkInitialized().authorize(toAuthzRequest(req)) + result != null && result.getDecision == AccessDecision.ALLOW } } diff --git a/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/SparkRangerAuditHandler.scala b/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/SparkRangerAuditHandler.scala deleted file mode 100644 index 5eaae852795..00000000000 --- a/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/ranger/SparkRangerAuditHandler.scala +++ /dev/null @@ -1,26 +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.kyuubi.plugin.spark.authz.ranger - -import org.apache.ranger.plugin.audit.RangerDefaultAuditHandler - -class SparkRangerAuditHandler extends RangerDefaultAuditHandler { - - // Implementing meaningfully audit functions - -} diff --git a/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/rule/rowfilter/FilteredShowObjectsExec.scala b/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/rule/rowfilter/FilteredShowObjectsExec.scala index fd617161bbb..455bf9e4897 100644 --- a/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/rule/rowfilter/FilteredShowObjectsExec.scala +++ b/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/rule/rowfilter/FilteredShowObjectsExec.scala @@ -51,8 +51,7 @@ object FilteredShowNamespaceExec extends FilteredShowObjectsCheck { val database = r.getString(0) val resource = AccessResource(ObjectType.DATABASE, database, null, null) val request = AccessRequest(resource, ugi, OperationType.SHOWDATABASES, AccessType.USE) - val result = SparkRangerAdminPlugin.isAccessAllowed(request) - result != null && result.getIsAllowed + SparkRangerAdminPlugin.isAccessAllowed(request) } } @@ -75,7 +74,6 @@ object FilteredShowTablesExec extends FilteredShowObjectsCheck { val objectType = if (isTemp) ObjectType.VIEW else ObjectType.TABLE val resource = AccessResource(objectType, database, table, null) val request = AccessRequest(resource, ugi, OperationType.SHOWTABLES, AccessType.USE) - val result = SparkRangerAdminPlugin.isAccessAllowed(request) - result != null && result.getIsAllowed + SparkRangerAdminPlugin.isAccessAllowed(request) } } diff --git a/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/rule/rowfilter/RuleReplaceShowObjectCommands.scala b/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/rule/rowfilter/RuleReplaceShowObjectCommands.scala index 0726a40e001..8cb9b8f9ca7 100644 --- a/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/rule/rowfilter/RuleReplaceShowObjectCommands.scala +++ b/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/rule/rowfilter/RuleReplaceShowObjectCommands.scala @@ -58,8 +58,7 @@ case class FilteredShowTablesCommand(delegated: RunnableCommand) val resource = AccessResource(objectType, database, table, null) val accessType = if (isExtended) AccessType.SELECT else AccessType.USE val request = AccessRequest(resource, ugi, OperationType.SHOWTABLES, accessType) - val result = SparkRangerAdminPlugin.isAccessAllowed(request) - result != null && result.getIsAllowed + SparkRangerAdminPlugin.isAccessAllowed(request) } } @@ -92,8 +91,7 @@ case class FilteredShowFunctionsCommand(delegated: RunnableCommand) val resource = AccessResource(ObjectType.FUNCTION, items(0), items(1), null) val request = AccessRequest(resource, ugi, OperationType.SHOWFUNCTIONS, AccessType.USE) - val result = SparkRangerAdminPlugin.isAccessAllowed(request) - result != null && result.getIsAllowed + SparkRangerAdminPlugin.isAccessAllowed(request) } } @@ -113,7 +111,6 @@ case class FilteredShowColumnsCommand(delegated: RunnableCommand) override protected def isAllowed(r: Row, ugi: UserGroupInformation): Boolean = { val resource = AccessResource(ObjectType.COLUMN, r.getString(0), r.getString(1), r.getString(2)) val request = AccessRequest(resource, ugi, OperationType.SHOWCOLUMNS, AccessType.USE) - val result = SparkRangerAdminPlugin.isAccessAllowed(request) - result != null && result.getIsAllowed + SparkRangerAdminPlugin.isAccessAllowed(request) } } diff --git a/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/util/AuthZUtils.scala b/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/util/AuthZUtils.scala index 1bac3642aeb..8ebd7d511b6 100644 --- a/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/util/AuthZUtils.scala +++ b/extensions/spark/kyuubi-spark-authz/src/main/scala/org/apache/kyuubi/plugin/spark/authz/util/AuthZUtils.scala @@ -25,7 +25,6 @@ import java.util.Base64 import org.apache.commons.lang3.StringUtils import org.apache.hadoop.security.UserGroupInformation -import org.apache.ranger.plugin.service.RangerBasePlugin import org.apache.spark.{SPARK_VERSION, SparkContext} import org.apache.spark.sql.SparkSession import org.apache.spark.sql.catalyst.plans.logical.{LogicalPlan, View} @@ -33,7 +32,6 @@ import org.apache.spark.sql.catalyst.plans.logical.{LogicalPlan, View} import org.apache.kyuubi.plugin.spark.authz.AccessControlException import org.apache.kyuubi.plugin.spark.authz.util.ReservedKeys._ import org.apache.kyuubi.util.SemanticVersion -import org.apache.kyuubi.util.reflect.DynConstructors import org.apache.kyuubi.util.reflect.ReflectUtils._ private[authz] object AuthZUtils { @@ -76,21 +74,6 @@ private[authz] object AuthZUtils { spark.conf.getOption(SKIP_CATALOGLESS_V2_RELATION_ENABLED_KEY) .exists(_.equalsIgnoreCase("true")) - lazy val isRanger21orGreater: Boolean = { - try { - DynConstructors.builder().impl( - classOf[RangerBasePlugin], - classOf[String], - classOf[String], - classOf[String]) - .buildChecked[RangerBasePlugin]() - true - } catch { - case _: NoSuchMethodException => - false - } - } - lazy val SPARK_RUNTIME_VERSION: SemanticVersion = SemanticVersion(SPARK_VERSION) lazy val isSparkV40OrGreater: Boolean = SPARK_RUNTIME_VERSION >= "4.0" diff --git a/extensions/spark/kyuubi-spark-authz/src/test/resources/ranger-spark-security.xml b/extensions/spark/kyuubi-spark-authz/src/test/resources/ranger-spark-security.xml index 337f8e1bc43..8b811c00ca9 100644 --- a/extensions/spark/kyuubi-spark-authz/src/test/resources/ranger-spark-security.xml +++ b/extensions/spark/kyuubi-spark-authz/src/test/resources/ranger-spark-security.xml @@ -16,6 +16,16 @@ ~ limitations under the License. --> + + ranger.authorizer.impl.class + org.apache.ranger.authz.embedded.RangerEmbeddedAuthorizer + + The tests exercise the embedded Ranger authorizer that pulls policies from + the local RangerLocalClient; the PDP remote authorizer is covered by + RangerRemoteAuthorizerSuite + + + ranger.plugin.spark.service.name hive_jenkins diff --git a/extensions/spark/kyuubi-spark-authz/src/test/scala/org/apache/kyuubi/plugin/spark/authz/ranger/AccessResourceSuite.scala b/extensions/spark/kyuubi-spark-authz/src/test/scala/org/apache/kyuubi/plugin/spark/authz/ranger/AccessResourceSuite.scala index 2de223faed1..3fa4b5f05a2 100644 --- a/extensions/spark/kyuubi-spark-authz/src/test/scala/org/apache/kyuubi/plugin/spark/authz/ranger/AccessResourceSuite.scala +++ b/extensions/spark/kyuubi-spark-authz/src/test/scala/org/apache/kyuubi/plugin/spark/authz/ranger/AccessResourceSuite.scala @@ -17,7 +17,12 @@ package org.apache.kyuubi.plugin.spark.authz.ranger +import scala.collection.JavaConverters._ + +import org.apache.ranger.authz.util.RangerResourceNameParser + import org.apache.kyuubi.KyuubiFunSuite +import org.apache.kyuubi.plugin.spark.authz.AccessControlException import org.apache.kyuubi.plugin.spark.authz.ObjectType._ class AccessResourceSuite extends KyuubiFunSuite { @@ -42,7 +47,7 @@ class AccessResourceSuite extends KyuubiFunSuite { assert(resource.catalog.isEmpty) assert(resource2.getDatabase === "my_db_name") assert(resource2.getTable === null) - assert(resource2.getValue("udf") === "my_func_name") + assert(resource2.getUdf === "my_func_name") assert(resource1.getColumn === null) assert(resource1.getColumns.isEmpty) @@ -79,4 +84,82 @@ class AccessResourceSuite extends KyuubiFunSuite { assert(resource1.getColumn === "my_col_1,my_col_2") assert(resource1.getColumns === Seq("my_col_1", "my_col_2")) } + + test("escape RRN metacharacters in resource names") { + assert( + AccessResource(TABLE, "my/db", "my\\table", null).toResourceInfos.head.getName === + s"table:my\\/db/my\\\\table") + + assert( + AccessResource(COLUMN, "d", "t", "c/1").toResourceInfos.head.getName === "column:d/t/c\\/1") + + val columns = AccessResource(COLUMN, "d", "t", "c1,c/2").toResourceInfos.head + assert(columns.getName === "column:d/t") + assert(columns.getSubResources.asScala.toSeq.sorted === + Seq("column:c1", "column:c\\/2")) + + assert(AccessResource( + FUNCTION, + "d", + "u/f", + null).toResourceInfos.head.getName === "udf:d/u\\/f") + + assert(AccessResource(URI, "/tmp/a b", null, null).toResourceInfos.map(_.getName) === + Seq(s"url:\\/tmp\\/a b", s"url:\\/tmp\\/a b/")) + } + + test("generate two uri resource variants as alternatives") { + assert(AccessResource(URI, "/tmp/data", null, null).toResourceInfos.map(_.getName) === + Seq(s"url:\\/tmp\\/data", s"url:\\/tmp\\/data/")) + } + + test("RRN metacharacter escaping round-trips through RangerResourceNameParser") { + def parse(name: String, template: String): java.util.Map[String, String] = + new RangerResourceNameParser(template).parseToMap(name.substring(name.indexOf(':') + 1)) + + val tableComponents = + parse( + AccessResource(TABLE, "my/db", "my\\table", null).toResourceInfos.head.getName, + "database/table") + assert(tableComponents.get("database") === "my/db") + assert(tableComponents.get("table") === "my\\table") + + val columnComponents = + parse( + AccessResource(COLUMN, "d", "t", "c/1").toResourceInfos.head.getName, + "database/table/column") + assert(columnComponents.get("column") === "c/1") + } + + test("deny blank resource names") { + intercept[AccessControlException](AccessResource(DATABASE, null, null, null).toResourceInfos) + intercept[AccessControlException](AccessResource(TABLE, null, "t", null).toResourceInfos) + intercept[AccessControlException](AccessResource(TABLE, "d", " ", null).toResourceInfos) + intercept[AccessControlException](AccessResource(COLUMN, null, "t", "c").toResourceInfos) + intercept[AccessControlException](AccessResource(COLUMN, "d", null, "c").toResourceInfos) + intercept[AccessControlException](AccessResource(FUNCTION, "d", null, null).toResourceInfos) + intercept[AccessControlException](AccessResource(URI, null, null, null).toResourceInfos) + } + + test("blank function database is checked against any database") { + assert(AccessResource(FUNCTION, "", "func", null).toResourceInfos.head.getName === "udf:*/func") + assert(AccessResource( + FUNCTION, + null, + "func", + null).toResourceInfos.head.getName === "udf:*/func") + } + + test("blank view database is checked against any database") { + // a local temporary view has no database in the SHOW TABLES output + assert(AccessResource(VIEW, "", "temp_view", null).toResourceInfos.head.getName === + "table:*/temp_view") + assert(AccessResource(VIEW, null, "temp_view", null).toResourceInfos.head.getName === + "table:*/temp_view") + assert(AccessResource( + VIEW, + "my_db_name", + "my_view_name", + null).toResourceInfos.head.getName === "table:my_db_name/my_view_name") + } } diff --git a/extensions/spark/kyuubi-spark-authz/src/test/scala/org/apache/kyuubi/plugin/spark/authz/ranger/MockPdpServer.scala b/extensions/spark/kyuubi-spark-authz/src/test/scala/org/apache/kyuubi/plugin/spark/authz/ranger/MockPdpServer.scala new file mode 100644 index 00000000000..e5cf8822363 --- /dev/null +++ b/extensions/spark/kyuubi-spark-authz/src/test/scala/org/apache/kyuubi/plugin/spark/authz/ranger/MockPdpServer.scala @@ -0,0 +1,120 @@ +/* + * 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.kyuubi.plugin.spark.authz.ranger + +import java.net.InetSocketAddress +import java.util.Properties +import java.util.concurrent.Executors + +import scala.collection.JavaConverters._ + +import com.fasterxml.jackson.databind.DeserializationFeature +import com.fasterxml.jackson.databind.json.JsonMapper +import com.sun.net.httpserver.{HttpExchange, HttpHandler, HttpServer} +import org.apache.hadoop.conf.Configuration +import org.apache.ranger.authz.api.RangerAuthorizer +import org.apache.ranger.authz.embedded.RangerEmbeddedAuthorizer +import org.apache.ranger.authz.model._ + +/** + * A test double of the Ranger PDP server, which authorizes requests with + * [[RangerEmbeddedAuthorizer]] backed by the local policy file, serving the same + * REST APIs as org.apache.ranger:pdp does. + */ +class MockPdpServer extends AutoCloseable { + + private val mapper = new JsonMapper() + .configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false) + + private val properties: Properties = { + // the embedded authorizer reads configurations from the given properties only, + // so feed the same ones the client-side plugin uses + val conf = new Configuration + conf.addResource(classOf[MockPdpServer].getClassLoader.getResource("ranger-spark-security.xml")) + val ret = new Properties + conf.iterator.asScala + .filter(entry => entry.getKey.startsWith("ranger.") || entry.getKey.startsWith("xasecure.")) + .foreach(entry => ret.put(entry.getKey, entry.getValue)) + ret.put("ranger.authz.init.services", "hive_jenkins") + ret.put("ranger.authz.service.hive_jenkins.servicetype", "spark") + ret.put("ranger.authz.app.type", "ranger-authz-pdp") + ret + } + + private val authorizer: RangerAuthorizer = new RangerEmbeddedAuthorizer(properties) + + private val authorizeHandler: HttpHandler = new HttpHandler { + override def handle(exchange: HttpExchange): Unit = { + val request = mapper.readValue(readRequestBody(exchange), classOf[RangerAuthzRequest]) + val result = authorizer.authorize(request) + respond(exchange, mapper.writeValueAsString(result)) + } + } + + private val authorizeMultiHandler: HttpHandler = new HttpHandler { + override def handle(exchange: HttpExchange): Unit = { + val request = mapper.readValue(readRequestBody(exchange), classOf[RangerMultiAuthzRequest]) + val result = authorizer.authorize(request) + respond(exchange, mapper.writeValueAsString(result)) + } + } + + private val executor = Executors.newFixedThreadPool(4) + + private val server: HttpServer = { + authorizer.init() + val ret = HttpServer.create(new InetSocketAddress("localhost", 0), 0) + ret.createContext("/authz/v1/authorize", authorizeHandler) + ret.createContext("/authz/v1/authorizeMulti", authorizeMultiHandler) + ret.setExecutor(executor) + ret.start() + ret + } + + val url: String = s"http://localhost:${server.getAddress.getPort}" + + private def readRequestBody(exchange: HttpExchange): Array[Byte] = { + val out = new java.io.ByteArrayOutputStream + val in = exchange.getRequestBody + val buf = new Array[Byte](8192) + var len = in.read(buf) + while (len != -1) { + out.write(buf, 0, len) + len = in.read(buf) + } + out.toByteArray + } + + private def respond(exchange: HttpExchange, body: String): Unit = { + val bytes = body.getBytes(java.nio.charset.StandardCharsets.UTF_8) + exchange.getResponseHeaders.add("Content-Type", "application/json") + exchange.sendResponseHeaders(200, bytes.length) + val out = exchange.getResponseBody + try out.write(bytes) + finally { + out.close() + exchange.close() + } + } + + override def close(): Unit = { + server.stop(0) + executor.shutdownNow() + authorizer.close() + } +} diff --git a/extensions/spark/kyuubi-spark-authz/src/test/scala/org/apache/kyuubi/plugin/spark/authz/ranger/RangerRemoteAuthorizerSuite.scala b/extensions/spark/kyuubi-spark-authz/src/test/scala/org/apache/kyuubi/plugin/spark/authz/ranger/RangerRemoteAuthorizerSuite.scala new file mode 100644 index 00000000000..5c09daf6458 --- /dev/null +++ b/extensions/spark/kyuubi-spark-authz/src/test/scala/org/apache/kyuubi/plugin/spark/authz/ranger/RangerRemoteAuthorizerSuite.scala @@ -0,0 +1,184 @@ +/* + * 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.kyuubi.plugin.spark.authz.ranger + +import java.util.Properties + +import scala.collection.JavaConverters._ + +import org.apache.hadoop.security.UserGroupInformation + +import org.apache.kyuubi.KyuubiFunSuite +import org.apache.kyuubi.plugin.spark.authz.{AccessControlException, OperationType} +import org.apache.kyuubi.plugin.spark.authz.ObjectType._ +import org.apache.kyuubi.plugin.spark.authz.RangerTestUsers._ +import org.apache.kyuubi.plugin.spark.authz.ranger.AccessType._ + +/** + * The tests exercise the Ranger PDP remote authorizer (the default), which sends + * authorization requests to a Ranger PDP server via REST APIs. + * + * By default, the requests go to [[MockPdpServer]] backed by the same policies used + * by the other suites. To validate against a real Ranger PDP server, start one + * configured with the same service (hive_jenkins) and policies, then run this suite + * with the Ranger PDP server url and the Ranger PDP client settings, e.g. the + * authentication settings, specified as system properties, e.g. + * {{{ + * mvn test -pl extensions/spark/kyuubi-spark-authz \ + * -DforkMode=never \ + * -DwildcardSuites=org.apache.kyuubi.plugin.spark.authz.ranger.RangerRemoteAuthorizerSuite \ + * -Dranger.authz.remote.pdp.url=https://ranger-pdp.org:8585 \ + * -Dranger.authz.remote.authn.type=header \ + * -Dranger.authz.remote.authn.header.X-Forwarded-User=ranger + * }}} + * + * The plugin reads the security-critical Ranger configurations from the configuration + * resources only, so the suite passes them to the plugin explicitly; the system + * properties above are consumed by this suite. + * + * Note that `-DforkMode=never` is required to run the tests in the maven JVM so + * that the system properties reach the suite. + */ +class RangerRemoteAuthorizerSuite extends KyuubiFunSuite { + + private val RemoteAuthorizer = "org.apache.ranger.authz.remote.RangerRemoteAuthorizer" + + // when a real Ranger PDP server is given via the system property, + // the mock server is not started + private lazy val mockPdp: Option[MockPdpServer] = + if (System.getProperty("ranger.authz.remote.pdp.url") == null) { + Some(new MockPdpServer) + } else { + None + } + + private def pdpUrl: String = + System.getProperty("ranger.authz.remote.pdp.url", mockPdp.map(_.url).orNull) + + private def ugiOf(user: String): UserGroupInformation = + UserGroupInformation.createRemoteUser(user) + + private def tableReq(user: String, db: String, table: String, accessType: AccessType) = + AccessRequest( + AccessResource(TABLE, db, table, null), + ugiOf(user), + OperationType.QUERY, + accessType) + + private def initializeRemoteAuthorizer(): Unit = { + // ensures the Ranger configuration resources are loaded + SparkRangerAdminPlugin.initialize() + val properties = new Properties + // the Ranger PDP client settings, e.g. the authentication settings, + // can be specified as system properties of the test JVM + System.getProperties.asScala + .filter { case (key, _) => key.startsWith("ranger.authz.remote.") } + .foreach { case (key, value) => properties.put(key, value) } + properties.setProperty("ranger.authorizer.impl.class", RemoteAuthorizer) + properties.setProperty("ranger.authz.remote.pdp.url", pdpUrl) + SparkRangerAdminPlugin.reset() + SparkRangerAdminPlugin.initialize(properties) + } + + override def afterAll(): Unit = { + SparkRangerAdminPlugin.reset() + SparkRangerAdminPlugin.initialize() + mockPdp.foreach(_.close()) + super.afterAll() + } + + // the authorizer may be re-initialized by other suites running in the same JVM + override def beforeEach(): Unit = { + initializeRemoteAuthorizer() + super.beforeEach() + } + + test("verify allowed access") { + SparkRangerAdminPlugin.verify(Seq(tableReq(bob, "default", "src", SELECT))) + SparkRangerAdminPlugin.verify( + Seq( + AccessRequest( + AccessResource(COLUMN, "default", "src", "key"), + ugiOf(kent), + OperationType.QUERY, + SELECT))) + } + + test("verify denied access") { + val e = intercept[AccessControlException] { + SparkRangerAdminPlugin.verify(Seq(tableReq(someone, "default", "src", SELECT))) + } + assert(e.getMessage.contains(s"does not have [select] privilege on [default/src]")) + } + + test("verify multiple denied accesses in single call") { + val requests = Seq( + tableReq(someone, "default", "src", SELECT), + tableReq(someone, "default", "perm_view", SELECT)) + val e = intercept[AccessControlException] { + SparkRangerAdminPlugin.verify(requests) + } + assert(e.getMessage.contains("[select] privilege on [default/src,default/perm_view]")) + } + + test("get filter expr") { + val filterExpr = + SparkRangerAdminPlugin.getFilterExpr(tableReq(bob, "default", "src", SELECT)) + assert(filterExpr.nonEmpty && filterExpr.get == "key<20") + } + + test("get mask expr") { + val maskedExpr = SparkRangerAdminPlugin.getMaskingExpr( + AccessRequest( + AccessResource(COLUMN, "default", "src", "value1"), + ugiOf(bob), + OperationType.QUERY, + SELECT)) + assert(maskedExpr.nonEmpty && maskedExpr.get == "md5(cast(value1 as string))") + + val showFirst4Expr = SparkRangerAdminPlugin.getMaskingExpr( + AccessRequest( + AccessResource(COLUMN, "default", "src", "value3"), + ugiOf(bob), + OperationType.QUERY, + SELECT)) + assert(showFirst4Expr.nonEmpty && showFirst4Expr.get.contains("regexp_replace")) + } + + test("access allowed for any access type") { + assert(SparkRangerAdminPlugin.isAccessAllowed( + tableReq(bob, "default_bob", "table_use1", USE))) + assert(!SparkRangerAdminPlugin.isAccessAllowed( + tableReq(someone, "default_bob", "table_use1", USE))) + } + + test("access allowed for uri") { + val uri = AccessRequest( + AccessResource(URI, "/tmp/data", null, null), + ugiOf(admin), + OperationType.LOAD, + READ) + assert(SparkRangerAdminPlugin.isAccessAllowed(uri)) + val denied = AccessRequest( + AccessResource(URI, "/tmp/data", null, null), + ugiOf(someone), + OperationType.LOAD, + READ) + assert(!SparkRangerAdminPlugin.isAccessAllowed(denied)) + } +} diff --git a/extensions/spark/kyuubi-spark-authz/src/test/scala/org/apache/kyuubi/plugin/spark/authz/ranger/RangerSparkExtensionSuite.scala b/extensions/spark/kyuubi-spark-authz/src/test/scala/org/apache/kyuubi/plugin/spark/authz/ranger/RangerSparkExtensionSuite.scala index 8805f0cd527..1fc13038e4e 100644 --- a/extensions/spark/kyuubi-spark-authz/src/test/scala/org/apache/kyuubi/plugin/spark/authz/ranger/RangerSparkExtensionSuite.scala +++ b/extensions/spark/kyuubi-spark-authz/src/test/scala/org/apache/kyuubi/plugin/spark/authz/ranger/RangerSparkExtensionSuite.scala @@ -139,10 +139,10 @@ abstract class RangerSparkExtensionSuite extends KyuubiFunSuite val singleCallConfig = s"ranger.plugin.${SparkRangerAdminPlugin.getServiceType}.authorize.in.single.call" try { - SparkRangerAdminPlugin.getRangerConf.setBoolean(singleCallConfig, true) + SparkRangerAdminPlugin.config.setBoolean(singleCallConfig, true) f } finally { - SparkRangerAdminPlugin.getRangerConf.setBoolean(singleCallConfig, false) + SparkRangerAdminPlugin.config.setBoolean(singleCallConfig, false) } } diff --git a/extensions/spark/kyuubi-spark-authz/src/test/scala/org/apache/kyuubi/plugin/spark/authz/ranger/SparkRangerAdminPluginSuite.scala b/extensions/spark/kyuubi-spark-authz/src/test/scala/org/apache/kyuubi/plugin/spark/authz/ranger/SparkRangerAdminPluginSuite.scala index 83ad3ba2385..c95873eebe7 100644 --- a/extensions/spark/kyuubi-spark-authz/src/test/scala/org/apache/kyuubi/plugin/spark/authz/ranger/SparkRangerAdminPluginSuite.scala +++ b/extensions/spark/kyuubi-spark-authz/src/test/scala/org/apache/kyuubi/plugin/spark/authz/ranger/SparkRangerAdminPluginSuite.scala @@ -17,16 +17,57 @@ package org.apache.kyuubi.plugin.spark.authz.ranger +import java.util.Properties + import org.apache.hadoop.security.UserGroupInformation +import org.apache.ranger.authz.api.RangerAuthorizer +import org.apache.ranger.authz.model.{RangerAuthzRequest, RangerAuthzResult, RangerMultiAuthzRequest, RangerMultiAuthzResult, RangerResourcePermissions, RangerResourcePermissionsRequest} import org.apache.kyuubi.KyuubiFunSuite -import org.apache.kyuubi.plugin.spark.authz.{ObjectType, OperationType} +import org.apache.kyuubi.plugin.spark.authz.{AccessControlException, ObjectType, OperationType} import org.apache.kyuubi.plugin.spark.authz.RangerTestNamespace._ import org.apache.kyuubi.plugin.spark.authz.RangerTestUsers._ import org.apache.kyuubi.plugin.spark.authz.ranger.SparkRangerAdminPlugin._ +/** + * An authorizer that returns no result for any authorization request. + */ +class NoResultAuthorizer(properties: Properties) extends RangerAuthorizer(properties) { + override def init(): Unit = () + override def close(): Unit = {} + override def authorize(request: RangerAuthzRequest): RangerAuthzResult = null + override def authorize(request: RangerMultiAuthzRequest): RangerMultiAuthzResult = null + override def getResourcePermissions( + request: RangerResourcePermissionsRequest): RangerResourcePermissions = null +} + class SparkRangerAdminPluginSuite extends KyuubiFunSuite { + private val RemoteAuthorizer = "org.apache.ranger.authz.remote.RangerRemoteAuthorizer" + + private def tableSelectRequest(user: String, database: String, table: String): AccessRequest = + AccessRequest( + AccessResource(ObjectType.TABLE, database, table, null), + UserGroupInformation.createRemoteUser(user), + OperationType.QUERY, + AccessType.SELECT) + + /** + * initializes the plugin with the given authorizer properties, + * and restores the default authorizer after running the body + */ + private def withAuthorizer(properties: Properties)(body: => Unit): Unit = { + SparkRangerAdminPlugin.initialize() + try { + SparkRangerAdminPlugin.reset() + SparkRangerAdminPlugin.initialize(properties) + body + } finally { + SparkRangerAdminPlugin.reset() + SparkRangerAdminPlugin.initialize() + } + } + test("get filter expression") { val bob = UserGroupInformation.createRemoteUser("bob") val are = AccessResource(ObjectType.TABLE, defaultDb, "src", null) @@ -65,4 +106,74 @@ class SparkRangerAdminPluginSuite extends KyuubiFunSuite { assert(maybeString.isEmpty) } } + + test("deny accesses when the authorizer returns no result") { + val properties = new Properties + properties.setProperty("ranger.authorizer.impl.class", classOf[NoResultAuthorizer].getName) + withAuthorizer(properties) { + val e = intercept[AccessControlException] { + verify(Seq(tableSelectRequest(bob, defaultDb, "src"))) + } + assert(e.getMessage.contains(s"Permission denied: user [$bob] does not have " + + s"[select] privilege on [$defaultDb/src]")) + } + } + + test("deny accesses when the authorizer returns no result in single call") { + val properties = new Properties + properties.setProperty("ranger.authorizer.impl.class", classOf[NoResultAuthorizer].getName) + withAuthorizer(properties) { + // the single call mode is read from the Ranger configuration resources + val config = SparkRangerAdminPlugin.config + config.setBoolean("ranger.plugin.spark.authorize.in.single.call", true) + try { + val e = intercept[AccessControlException] { + verify(Seq( + tableSelectRequest(bob, defaultDb, "src"), + tableSelectRequest(bob, defaultDb, "perm_view"))) + } + assert(e.getMessage.contains(s"Permission denied: user [$bob] does not have " + + s"[select] privilege on [$defaultDb/src,$defaultDb/perm_view]")) + } finally { + config.unset("ranger.plugin.spark.authorize.in.single.call") + } + } + } + + test("the security-critical configurations cannot be overridden by system properties") { + System.setProperty("ranger.authorizer.impl.class", RemoteAuthorizer) + // an unreachable Ranger PDP server, to detect an unwanted switch to the remote + // authorizer by a failed access check + System.setProperty("ranger.authz.remote.pdp.url", "http://localhost:1") + try { + SparkRangerAdminPlugin.reset() + SparkRangerAdminPlugin.initialize() + // the embedded authorizer pinned in the configuration resources is used, + // so the access is evaluated against the local policies + assert(isAccessAllowed(tableSelectRequest(admin, defaultDb, "src"))) + } finally { + System.clearProperty("ranger.authorizer.impl.class") + System.clearProperty("ranger.authz.remote.pdp.url") + SparkRangerAdminPlugin.reset() + SparkRangerAdminPlugin.initialize() + } + } + + test("fail to initialize when the Ranger PDP url is missing for the remote authorizer") { + SparkRangerAdminPlugin.initialize() + val config = SparkRangerAdminPlugin.config + val originalImplClass = config.get("ranger.authorizer.impl.class") + config.set("ranger.authorizer.impl.class", RemoteAuthorizer) + try { + SparkRangerAdminPlugin.reset() + val e = intercept[IllegalArgumentException] { + SparkRangerAdminPlugin.initialize() + } + assert(e.getMessage.contains("ranger.authz.remote.pdp.url")) + } finally { + config.set("ranger.authorizer.impl.class", originalImplClass) + SparkRangerAdminPlugin.reset() + SparkRangerAdminPlugin.initialize() + } + } }