diff --git a/CHANGES.txt b/CHANGES.txt index d80eeaf168..30a741e718 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 3.0 + * Add role based access control (CASSANDRA-7653) * Group sstables for anticompaction correctly (CASSANDRA-8578) * Add ReadFailureException to native protocol, respond immediately when replicas encounter errors while handling diff --git a/NEWS.txt b/NEWS.txt index 8d8ebdcf3d..b9c41734ba 100644 --- a/NEWS.txt +++ b/NEWS.txt @@ -18,6 +18,14 @@ using the provided 'sstableupgrade' tool. New features ------------ + - Authentication & Authorization APIs have been updated to introduce + roles. Roles and Permissions granted to them are inherited, supporting + role based access control. The role concept supercedes that of users + and CQL constructs such as CREATE USER are deprecated but retained for + compatibility. The requirement to explicitly create Roles in Cassandra + even when auth is handled by an external system has been removed, so + authentication & authorization can be delegated to such systems in their + entirety. - SSTable file name is changed. Now you don't have Keyspace/CF name in file name. Also, secondary index has its own directory under parent's directory. @@ -25,6 +33,18 @@ New features Upgrading --------- + - IAuthenticator been updated to remove responsibility for user/role + maintenance and is now solely responsible for validating credentials, + This is primarily done via SASL, though an optional method exists for + systems which need support for the Thrift login() method. + - IRoleManager interface has been added which takes over the maintenance + functions from IAuthenticator. IAuthorizer is mainly unchanged. Auth data + in systems using the stock internal implementations PasswordAuthenticator + & CassandraAuthorizer will be automatically converted during upgrade, + with minimal operator intervention required. Custom implementations will + require modification, though these can be used in conjunction with the + stock CassandraRoleManager so providing an IRoleManager implementation + should not usually be necessary. - Fat client support has been removed since we have push notifications to clients - cassandra-cli has been removed. Please use cqlsh instead. - YamlFileNetworkTopologySnitch has been removed; switch to @@ -37,6 +57,7 @@ Upgrading in the normal order and not anymore in the order in which the column values were specified in the IN restriction. + 2.1.2 ===== diff --git a/bin/cqlsh b/bin/cqlsh index 363a4f64de..427dcac17d 100755 --- a/bin/cqlsh +++ b/bin/cqlsh @@ -108,7 +108,7 @@ except ImportError, e: from cassandra.cluster import Cluster, PagedResult from cassandra.query import SimpleStatement, ordered_dict_factory from cassandra.policies import WhiteListRoundRobinPolicy -from cassandra.metadata import protect_name, protect_names, protect_value +from cassandra.metadata import protect_name, protect_names, protect_value, KeyspaceMetadata, TableMetadata, ColumnMetadata from cassandra.auth import PlainTextAuthProvider # cqlsh should run correctly when run out of a Cassandra source tree, @@ -751,9 +751,31 @@ class Shell(cmd.Cmd): ksmeta = self.get_keyspace_meta(ksname) if tablename not in ksmeta.tables: - raise ColumnFamilyNotFound("Column family %r not found" % tablename) + if ksname == 'system_auth' and tablename in ['roles','role_permissions']: + self.get_fake_auth_table_meta(ksname, tablename) + else: + raise ColumnFamilyNotFound("Column family %r not found" % tablename) + else: + return ksmeta.tables[tablename] - return ksmeta.tables[tablename] + def get_fake_auth_table_meta(self, ksname, tablename): + # may be using external auth implementation so internal tables + # aren't actually defined in schema. In this case, we'll fake + # them up + if tablename == 'roles': + ks_meta = KeyspaceMetadata(ksname, True, None, None) + table_meta = TableMetadata(ks_meta, 'roles') + table_meta.columns['role'] = ColumnMetadata(table_meta, 'role', cassandra.cqltypes.UTF8Type) + table_meta.columns['is_superuser'] = ColumnMetadata(table_meta, 'is_superuser', cassandra.cqltypes.BooleanType) + table_meta.columns['can_login'] = ColumnMetadata(table_meta, 'can_login', cassandra.cqltypes.BooleanType) + elif tablename == 'role_permissions': + ks_meta = KeyspaceMetadata(ksname, True, None, None) + table_meta = TableMetadata(ks_meta, 'role_permissions') + table_meta.columns['role'] = ColumnMetadata(table_meta, 'role', cassandra.cqltypes.UTF8Type) + table_meta.columns['resource'] = ColumnMetadata(table_meta, 'resource', cassandra.cqltypes.UTF8Type) + table_meta.columns['permission'] = ColumnMetadata(table_meta, 'permission', cassandra.cqltypes.UTF8Type) + else: + raise ColumnFamilyNotFoundException("Column family %r not found" % tablename) def get_usertypes_meta(self): data = self.session.execute("select * from system.schema_usertypes") @@ -1006,10 +1028,10 @@ class Shell(cmd.Cmd): if statement.query_string[:6].lower() == 'select': self.print_result(rows, self.parse_for_table_meta(statement.query_string)) - elif statement.query_string.lower().startswith("list users"): - self.print_result(rows, self.get_table_meta('system_auth','users')) + elif statement.query_string.lower().startswith("list users") or statement.query_string.lower().startswith("list roles"): + self.print_result(rows, self.get_table_meta('system_auth','roles')) elif statement.query_string.lower().startswith("list"): - self.print_result(rows, self.get_table_meta('system_auth','permissions')) + self.print_result(rows, self.get_table_meta('system_auth','role_permissions')) elif rows: # CAS INSERT/UPDATE self.writeresult("") diff --git a/conf/cassandra.yaml b/conf/cassandra.yaml index ca2ca1bb09..24bab099d5 100644 --- a/conf/cassandra.yaml +++ b/conf/cassandra.yaml @@ -62,6 +62,7 @@ batchlog_replay_throttle_in_kb: 1024 # - PasswordAuthenticator relies on username/password pairs to authenticate # users. It keeps usernames and hashed passwords in system_auth.credentials table. # Please increase system_auth keyspace replication factor if you use this authenticator. +# If using PasswordAuthenticator, CassandraRoleManager must also be used (see below) authenticator: AllowAllAuthenticator # Authorization backend, implementing IAuthorizer; used to limit access/provide permissions @@ -73,6 +74,25 @@ authenticator: AllowAllAuthenticator # increase system_auth keyspace replication factor if you use this authorizer. authorizer: AllowAllAuthorizer +# Part of the Authentication & Authorization backend, implementing IRoleManager; used +# to maintain grants and memberships between roles. +# Out of the box, Cassandra provides org.apache.cassandra.auth.CassandraRoleManager, +# which stores role information in the system_auth keyspace. Most functions of the +# IRoleManager require an authenticated login, so unless the configured IAuthenticator +# actually implements authentication, most of this functionality will be unavailable. +# +# - CassandraRoleManager stores role data in the system_auth keyspace. Please +# increase system_auth keyspace replication factor if you use this role manager. +role_manager: CassandraRoleManager + +# Validity period for roles cache (fetching permissions can be an +# expensive operation depending on the authorizer). Granted roles are cached for +# authenticated sessions in AuthenticatedUser and after the period specified +# here, become eligible for (async) reload. +# Defaults to 2000, set to 0 to disable. +# Will be disabled automatically for AllowAllAuthenticator. +roles_validity_in_ms: 2000 + # Validity period for permissions cache (fetching permissions can be an # expensive operation depending on the authorizer, CassandraAuthorizer is # one example). Defaults to 2000, set to 0 to disable. diff --git a/pylib/cqlshlib/cql3handling.py b/pylib/cqlshlib/cql3handling.py index a376c33736..552f5c1ebe 100644 --- a/pylib/cqlshlib/cql3handling.py +++ b/pylib/cqlshlib/cql3handling.py @@ -254,10 +254,16 @@ JUNK ::= /([ \t\r\f\v]+|(--|[/][/])[^\n\r]*([\n\r]|$)|[/][*].*?[*][/])/ ; | | | + | + | + | + | ; ::= + | | + | | ; @@ -1169,14 +1175,49 @@ syntax_rules += r''' ''' syntax_rules += r''' - ::= "GRANT" "ON" "TO" + ::= + | + | + ; + + ::= "CREATE" "ROLE" + ( "WITH" ("AND" )*)? + ( "SUPERUSER" | "NOSUPERUSER" )? + ( "LOGIN" | "NOLOGIN" )? + ; + + ::= "ALTER" "ROLE" + ( "WITH" ("AND" )*)? + ( "SUPERUSER" | "NOSUPERUSER" )? + ( "LOGIN" | "NOLOGIN" )? + ; + ::= "PASSWORD" + | "OPTIONS" + ; + + ::= "DROP" "ROLE" + ; + + ::= "GRANT" "TO" + ; + + ::= "REVOKE" "FROM" + ; + + ::= "LIST" "ROLES" + ( "OF" )? "NORECURSIVE"? + ; +''' + +syntax_rules += r''' + ::= "GRANT" "ON" "TO" ; - ::= "REVOKE" "ON" "FROM" + ::= "REVOKE" "ON" "FROM" ; ::= "LIST" - ( "ON" )? ( "OF" )? "NORECURSIVE"? + ( "ON" )? ( "OF" )? "NORECURSIVE"? ; ::= "AUTHORIZE" @@ -1214,11 +1255,24 @@ def username_name_completer(ctxt, cass): session = cass.session return [maybe_quote(row.values()[0].replace("'", "''")) for row in session.execute("LIST USERS")] +@completer_for('rolename', 'role') +def rolename_completer(ctxt, cass): + def maybe_quote(name): + if CqlRuleSet.is_valid_cql3_name(name): + return name + return "'%s'" % name + + # disable completion for CREATE ROLE. + if ctxt.matched[0][0] == 'K_CREATE': + return [Hint('')] + + session = cass.session + return [maybe_quote(row[0].replace("'", "''")) for row in session.execute("LIST ROLES")] + syntax_rules += r''' ::= "CREATE" "TRIGGER" ( "IF" "NOT" "EXISTS" )? "ON" cf= "USING" class= ; - ::= "DROP" "TRIGGER" ( "IF" "EXISTS" )? triggername= "ON" cf= ; diff --git a/pylib/cqlshlib/helptopics.py b/pylib/cqlshlib/helptopics.py index cbc6597f3a..c4f65b0f66 100644 --- a/pylib/cqlshlib/helptopics.py +++ b/pylib/cqlshlib/helptopics.py @@ -666,7 +666,9 @@ class CQL3HelpTopics(CQLHelpTopics): def help_create(self): super(CQL3HelpTopics, self).help_create() - print " HELP CREATE_USER;\n" + print """ HELP CREATE_USER; + HELP CREATE_ROLE; + """ def help_alter(self): print """ @@ -702,8 +704,10 @@ class CQL3HelpTopics(CQLHelpTopics): """ def help_drop(self): - super(CQL3HelpTopics, self).help_drop() - print " HELP DROP_USER;\n" + super(CQL3HelpTopics, self).help_create() + print """ HELP DROP_USER; + HELP DROP_ROLE; + """ def help_list(self): print """ @@ -758,10 +762,10 @@ class CQL3HelpTopics(CQLHelpTopics): ON ALL KEYSPACES | KEYSPACE | [TABLE] [.] - TO + TO [ROLE | USER ] Grant the specified permission (or all permissions) on a resource - to a user. + to a role or user. To be able to grant a permission on some resource you have to have that permission yourself and also AUTHORIZE permission on it, @@ -776,10 +780,10 @@ class CQL3HelpTopics(CQLHelpTopics): ON ALL KEYSPACES | KEYSPACE | [TABLE] [.]
- FROM + FROM [ROLE | USER ] Revokes the specified permission (or all permissions) on a resource - from a user. + from a role or user. To be able to revoke a permission on some resource you have to have that permission yourself and also AUTHORIZE permission on it, @@ -794,12 +798,13 @@ class CQL3HelpTopics(CQLHelpTopics): [ON ALL KEYSPACES | KEYSPACE | [TABLE] [.]
] - [OF ] + [OF [ROLE | USER ] [NORECURSIVE] Omitting ON part will list permissions on ALL KEYSPACES, every keyspace and table. - Omitting OF part will list permissions of all users. + Omitting OF [ROLE | USER ] part will list permissions + of all roles and users. Omitting NORECURSIVE specifier will list permissions of the resource and all its parents (table, table's keyspace and ALL KEYSPACES). @@ -818,3 +823,46 @@ class CQL3HelpTopics(CQLHelpTopics): MODIFY: required for INSERT, DELETE, UPDATE, TRUNCATE SELECT: required for SELECT """ + + def help_create_role(self): + print """ + CREATE ROLE ; + + CREATE ROLE creates a new Cassandra role. + Only superusers can issue CREATE ROLE requests. + To create a superuser account use SUPERUSER option (NOSUPERUSER is the default). + """ + + def help_drop_role(self): + print """ + DROP ROLE ; + + DROP ROLE removes an existing role. You have to be logged in as a superuser + to issue a DROP ROLE statement. + """ + + def help_list_roles(self): + print """ + LIST ROLES [OF [ROLE | USER ] [NORECURSIVE]]; + + Only superusers can use the OF clause to list the roles granted to a role or user. + If a superuser omits the OF clause then all the created roles will be listed. + If a non-superuser calls LIST ROLES then the roles granted to that user are listed. + If NORECURSIVE is provided then only directly granted roles are listed. + """ + + def help_grant_role(self): + print """ + GRANT ROLE TO [ROLE | USER ] + + Grant the specified role to another role or user. You have to be logged + in as superuser to issue a GRANT ROLE statement. + """ + + def help_revoke_role(self): + print """ + REVOKE ROLE FROM [ROLE | USER ] + + Revoke the specified role from another role or user. You have to be logged + in as superuser to issue a REVOKE ROLE statement. + """ diff --git a/src/java/org/apache/cassandra/auth/AllowAllAuthenticator.java b/src/java/org/apache/cassandra/auth/AllowAllAuthenticator.java index def6045013..bc00c3e72c 100644 --- a/src/java/org/apache/cassandra/auth/AllowAllAuthenticator.java +++ b/src/java/org/apache/cassandra/auth/AllowAllAuthenticator.java @@ -23,45 +23,16 @@ import java.util.Set; import org.apache.cassandra.exceptions.AuthenticationException; import org.apache.cassandra.exceptions.ConfigurationException; -import org.apache.cassandra.exceptions.InvalidRequestException; public class AllowAllAuthenticator implements IAuthenticator { + private static final SaslNegotiator AUTHENTICATOR_INSTANCE = new Negotiator(); + public boolean requireAuthentication() { return false; } - public Set
", keyspace, columnFamily); + case TABLE: + return String.format("
", keyspace, table); } throw new AssertionError(); } @@ -239,12 +239,12 @@ public class DataResource implements IResource return Objects.equal(level, ds.level) && Objects.equal(keyspace, ds.keyspace) - && Objects.equal(columnFamily, ds.columnFamily); + && Objects.equal(table, ds.table); } @Override public int hashCode() { - return Objects.hashCode(level, keyspace, columnFamily); + return Objects.hashCode(level, keyspace, table); } } diff --git a/src/java/org/apache/cassandra/auth/IAuthenticator.java b/src/java/org/apache/cassandra/auth/IAuthenticator.java index 608649078e..24792f63b6 100644 --- a/src/java/org/apache/cassandra/auth/IAuthenticator.java +++ b/src/java/org/apache/cassandra/auth/IAuthenticator.java @@ -22,85 +22,15 @@ import java.util.Set; import org.apache.cassandra.exceptions.AuthenticationException; import org.apache.cassandra.exceptions.ConfigurationException; -import org.apache.cassandra.exceptions.RequestExecutionException; -import org.apache.cassandra.exceptions.RequestValidationException; public interface IAuthenticator { - static final String USERNAME_KEY = "username"; - static final String PASSWORD_KEY = "password"; - - /** - * Supported CREATE USER/ALTER USER options. - * Currently only PASSWORD is available. - */ - enum Option - { - PASSWORD - } - /** * Whether or not the authenticator requires explicit login. * If false will instantiate user with AuthenticatedUser.ANONYMOUS_USER. */ boolean requireAuthentication(); - /** - * Set of options supported by CREATE USER and ALTER USER queries. - * Should never return null - always return an empty set instead. - */ - Set
, we need to correct the resource. resource = maybeCorrectResource(resource, state); if (!resource.exists()) - throw new InvalidRequestException(String.format("%s doesn't exist", resource)); + throw new InvalidRequestException(String.format("Resource %s doesn't exist", resource)); } public void checkAccess(ClientState state) throws UnauthorizedException diff --git a/src/java/org/apache/cassandra/cql3/statements/RevokeRoleStatement.java b/src/java/org/apache/cassandra/cql3/statements/RevokeRoleStatement.java new file mode 100644 index 0000000000..98c2b4e867 --- /dev/null +++ b/src/java/org/apache/cassandra/cql3/statements/RevokeRoleStatement.java @@ -0,0 +1,40 @@ +/* + * 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.cassandra.cql3.statements; + +import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.cql3.RoleName; +import org.apache.cassandra.exceptions.RequestExecutionException; +import org.apache.cassandra.exceptions.RequestValidationException; +import org.apache.cassandra.service.ClientState; +import org.apache.cassandra.transport.messages.ResultMessage; + +public class RevokeRoleStatement extends RoleManagementStatement +{ + public RevokeRoleStatement(RoleName name, RoleName grantee) + { + super(name, grantee); + } + + public ResultMessage execute(ClientState state) throws RequestValidationException, RequestExecutionException + { + DatabaseDescriptor.getRoleManager().revokeRole(state.getUser(), role, grantee); + return null; + } + +} diff --git a/src/java/org/apache/cassandra/cql3/statements/RevokeStatement.java b/src/java/org/apache/cassandra/cql3/statements/RevokeStatement.java index 6f8ccd1cc0..7ce52599bf 100644 --- a/src/java/org/apache/cassandra/cql3/statements/RevokeStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/RevokeStatement.java @@ -1,4 +1,4 @@ -/** +/* * 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 @@ -7,14 +7,13 @@ * "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 + * 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. + * 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.cassandra.cql3.statements; @@ -23,6 +22,7 @@ import java.util.Set; import org.apache.cassandra.auth.DataResource; import org.apache.cassandra.auth.Permission; import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.cql3.RoleName; import org.apache.cassandra.exceptions.RequestExecutionException; import org.apache.cassandra.exceptions.RequestValidationException; import org.apache.cassandra.service.ClientState; @@ -30,14 +30,14 @@ import org.apache.cassandra.transport.messages.ResultMessage; public class RevokeStatement extends PermissionAlteringStatement { - public RevokeStatement(Set permissions, DataResource resource, String username) + public RevokeStatement(Set permissions, DataResource resource, RoleName grantee) { - super(permissions, resource, username); + super(permissions, resource, grantee); } public ResultMessage execute(ClientState state) throws RequestValidationException, RequestExecutionException { - DatabaseDescriptor.getAuthorizer().revoke(state.getUser(), permissions, resource, username); + DatabaseDescriptor.getAuthorizer().revoke(state.getUser(), permissions, resource, grantee); return null; } } diff --git a/src/java/org/apache/cassandra/cql3/statements/RoleManagementStatement.java b/src/java/org/apache/cassandra/cql3/statements/RoleManagementStatement.java new file mode 100644 index 0000000000..d67b42ce77 --- /dev/null +++ b/src/java/org/apache/cassandra/cql3/statements/RoleManagementStatement.java @@ -0,0 +1,54 @@ +/* + * 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.cassandra.cql3.statements; + +import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.cql3.RoleName; +import org.apache.cassandra.exceptions.InvalidRequestException; +import org.apache.cassandra.exceptions.RequestValidationException; +import org.apache.cassandra.exceptions.UnauthorizedException; +import org.apache.cassandra.service.ClientState; + +public abstract class RoleManagementStatement extends AuthorizationStatement +{ + protected final String role; + protected final String grantee; + + public RoleManagementStatement(RoleName name, RoleName grantee) + { + this.role = name.getName(); + this.grantee = grantee.getName(); + } + + public void checkAccess(ClientState state) throws UnauthorizedException, InvalidRequestException + { + if (!state.getUser().isSuper()) + throw new UnauthorizedException("Only superusers are allowed to perform role management queries"); + } + + public void validate(ClientState state) throws RequestValidationException + { + state.ensureNotAnonymous(); + + if (!DatabaseDescriptor.getRoleManager().isExistingRole(role)) + throw new InvalidRequestException(String.format("%s doesn't exist", role)); + + if (!DatabaseDescriptor.getRoleManager().isExistingRole(grantee)) + throw new InvalidRequestException(String.format("%s doesn't exist", grantee)); + } +} diff --git a/src/java/org/apache/cassandra/hadoop/AbstractBulkRecordWriter.java b/src/java/org/apache/cassandra/hadoop/AbstractBulkRecordWriter.java index 136c8dc2fe..5ba0a9696e 100644 --- a/src/java/org/apache/cassandra/hadoop/AbstractBulkRecordWriter.java +++ b/src/java/org/apache/cassandra/hadoop/AbstractBulkRecordWriter.java @@ -21,18 +21,13 @@ import java.io.Closeable; import java.io.IOException; import java.net.InetAddress; import java.net.UnknownHostException; -import java.util.HashMap; -import java.util.HashSet; -import java.util.Iterator; -import java.util.List; -import java.util.Map; -import java.util.Set; -import java.util.concurrent.ExecutionException; -import java.util.concurrent.Future; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.TimeoutException; +import java.util.*; +import java.util.concurrent.*; -import org.apache.cassandra.auth.IAuthenticator; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import org.apache.cassandra.auth.PasswordAuthenticator; import org.apache.cassandra.config.CFMetaData; import org.apache.cassandra.config.Config; import org.apache.cassandra.config.DatabaseDescriptor; @@ -46,8 +41,6 @@ import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.mapreduce.RecordWriter; import org.apache.hadoop.mapreduce.TaskAttemptContext; import org.apache.hadoop.util.Progressable; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; public abstract class AbstractBulkRecordWriter extends RecordWriter implements org.apache.hadoop.mapred.RecordWriter @@ -191,8 +184,8 @@ implements org.apache.hadoop.mapred.RecordWriter if (username != null) { Map creds = new HashMap(); - creds.put(IAuthenticator.USERNAME_KEY, username); - creds.put(IAuthenticator.PASSWORD_KEY, password); + creds.put(PasswordAuthenticator.USERNAME_KEY, username); + creds.put(PasswordAuthenticator.PASSWORD_KEY, password); AuthenticationRequest authRequest = new AuthenticationRequest(creds); client.login(authRequest); } diff --git a/src/java/org/apache/cassandra/hadoop/AbstractColumnFamilyInputFormat.java b/src/java/org/apache/cassandra/hadoop/AbstractColumnFamilyInputFormat.java index f4ad40f7d3..6fe2239e3c 100644 --- a/src/java/org/apache/cassandra/hadoop/AbstractColumnFamilyInputFormat.java +++ b/src/java/org/apache/cassandra/hadoop/AbstractColumnFamilyInputFormat.java @@ -20,36 +20,22 @@ package org.apache.cassandra.hadoop; import java.io.IOException; import java.net.InetAddress; import java.util.*; -import java.util.concurrent.Callable; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Future; -import java.util.concurrent.LinkedBlockingQueue; -import java.util.concurrent.ThreadPoolExecutor; -import java.util.concurrent.TimeUnit; +import java.util.concurrent.*; import com.google.common.collect.ImmutableList; import com.google.common.collect.Lists; +import org.apache.commons.lang3.StringUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.apache.cassandra.auth.IAuthenticator; +import org.apache.cassandra.auth.PasswordAuthenticator; import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.dht.Range; import org.apache.cassandra.dht.Token; -import org.apache.cassandra.thrift.AuthenticationRequest; -import org.apache.cassandra.thrift.Cassandra; -import org.apache.cassandra.thrift.CfSplit; -import org.apache.cassandra.thrift.InvalidRequestException; -import org.apache.cassandra.thrift.KeyRange; -import org.apache.cassandra.thrift.TokenRange; -import org.apache.commons.lang3.StringUtils; +import org.apache.cassandra.thrift.*; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.mapred.JobConf; -import org.apache.hadoop.mapreduce.InputFormat; -import org.apache.hadoop.mapreduce.InputSplit; -import org.apache.hadoop.mapreduce.JobContext; -import org.apache.hadoop.mapreduce.TaskAttemptContext; -import org.apache.hadoop.mapreduce.TaskAttemptID; +import org.apache.hadoop.mapreduce.*; import org.apache.thrift.TApplicationException; import org.apache.thrift.TException; import org.apache.thrift.protocol.TBinaryProtocol; @@ -106,8 +92,8 @@ public abstract class AbstractColumnFamilyInputFormat extends InputFormat< if ((ConfigHelper.getInputKeyspaceUserName(conf) != null) && (ConfigHelper.getInputKeyspacePassword(conf) != null)) { Map creds = new HashMap(); - creds.put(IAuthenticator.USERNAME_KEY, ConfigHelper.getInputKeyspaceUserName(conf)); - creds.put(IAuthenticator.PASSWORD_KEY, ConfigHelper.getInputKeyspacePassword(conf)); + creds.put(PasswordAuthenticator.USERNAME_KEY, ConfigHelper.getInputKeyspaceUserName(conf)); + creds.put(PasswordAuthenticator.PASSWORD_KEY, ConfigHelper.getInputKeyspacePassword(conf)); AuthenticationRequest authRequest = new AuthenticationRequest(creds); client.login(authRequest); } diff --git a/src/java/org/apache/cassandra/hadoop/AbstractColumnFamilyOutputFormat.java b/src/java/org/apache/cassandra/hadoop/AbstractColumnFamilyOutputFormat.java index f574641a71..03d00456f4 100644 --- a/src/java/org/apache/cassandra/hadoop/AbstractColumnFamilyOutputFormat.java +++ b/src/java/org/apache/cassandra/hadoop/AbstractColumnFamilyOutputFormat.java @@ -25,8 +25,9 @@ import java.util.Map; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.apache.cassandra.auth.IAuthenticator; -import org.apache.cassandra.thrift.*; +import org.apache.cassandra.auth.PasswordAuthenticator; +import org.apache.cassandra.thrift.AuthenticationRequest; +import org.apache.cassandra.thrift.Cassandra; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.mapreduce.*; import org.apache.thrift.protocol.TBinaryProtocol; @@ -134,8 +135,8 @@ public abstract class AbstractColumnFamilyOutputFormat extends OutputForma public static void login(String user, String password, Cassandra.Client client) throws Exception { Map creds = new HashMap(); - creds.put(IAuthenticator.USERNAME_KEY, user); - creds.put(IAuthenticator.PASSWORD_KEY, password); + creds.put(PasswordAuthenticator.USERNAME_KEY, user); + creds.put(PasswordAuthenticator.PASSWORD_KEY, password); AuthenticationRequest authRequest = new AuthenticationRequest(creds); client.login(authRequest); } diff --git a/src/java/org/apache/cassandra/hadoop/pig/AbstractCassandraStorage.java b/src/java/org/apache/cassandra/hadoop/pig/AbstractCassandraStorage.java index 0ffd4429e3..447c8cea00 100644 --- a/src/java/org/apache/cassandra/hadoop/pig/AbstractCassandraStorage.java +++ b/src/java/org/apache/cassandra/hadoop/pig/AbstractCassandraStorage.java @@ -25,27 +25,28 @@ import java.nio.ByteBuffer; import java.nio.charset.CharacterCodingException; import java.util.*; -import org.apache.cassandra.db.Cell; -import org.apache.cassandra.schema.LegacySchemaTables; -import org.apache.cassandra.db.SystemKeyspace; -import org.apache.cassandra.exceptions.ConfigurationException; -import org.apache.cassandra.exceptions.SyntaxException; -import org.apache.cassandra.auth.IAuthenticator; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import org.apache.cassandra.auth.PasswordAuthenticator; import org.apache.cassandra.config.CFMetaData; import org.apache.cassandra.config.ColumnDefinition; +import org.apache.cassandra.db.Cell; +import org.apache.cassandra.db.SystemKeyspace; import org.apache.cassandra.db.marshal.*; import org.apache.cassandra.db.marshal.AbstractCompositeType.CompositeComponent; +import org.apache.cassandra.exceptions.ConfigurationException; +import org.apache.cassandra.exceptions.SyntaxException; +import org.apache.cassandra.hadoop.ConfigHelper; +import org.apache.cassandra.schema.LegacySchemaTables; import org.apache.cassandra.serializers.CollectionSerializer; -import org.apache.cassandra.hadoop.*; import org.apache.cassandra.thrift.*; -import org.apache.cassandra.utils.ByteBufferUtil; -import org.apache.cassandra.utils.FBUtilities; -import org.apache.cassandra.utils.Hex; -import org.apache.cassandra.utils.UUIDGen; - +import org.apache.cassandra.utils.*; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; -import org.apache.hadoop.mapreduce.*; +import org.apache.hadoop.mapreduce.InputFormat; +import org.apache.hadoop.mapreduce.Job; +import org.apache.hadoop.mapreduce.OutputFormat; import org.apache.pig.*; import org.apache.pig.backend.executionengine.ExecException; import org.apache.pig.data.*; @@ -54,8 +55,6 @@ import org.apache.thrift.TDeserializer; import org.apache.thrift.TException; import org.apache.thrift.TSerializer; import org.apache.thrift.protocol.TBinaryProtocol; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; /** * A LoadStoreFunc for retrieving data from and storing data to Cassandra @@ -505,8 +504,8 @@ public abstract class AbstractCassandraStorage extends LoadFunc implements Store if (username != null && password != null) { Map credentials = new HashMap(2); - credentials.put(IAuthenticator.USERNAME_KEY, username); - credentials.put(IAuthenticator.PASSWORD_KEY, password); + credentials.put(PasswordAuthenticator.USERNAME_KEY, username); + credentials.put(PasswordAuthenticator.PASSWORD_KEY, password); try { diff --git a/src/java/org/apache/cassandra/service/ClientState.java b/src/java/org/apache/cassandra/service/ClientState.java index 36f8326c7c..21d10f9e17 100644 --- a/src/java/org/apache/cassandra/service/ClientState.java +++ b/src/java/org/apache/cassandra/service/ClientState.java @@ -18,11 +18,12 @@ package org.apache.cassandra.service; import java.net.SocketAddress; -import java.util.*; +import java.util.Arrays; +import java.util.HashSet; +import java.util.Set; import java.util.concurrent.atomic.AtomicLong; import com.google.common.collect.Iterables; -import com.google.common.collect.Sets; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -31,13 +32,13 @@ import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.config.Schema; import org.apache.cassandra.cql3.QueryHandler; import org.apache.cassandra.cql3.QueryProcessor; -import org.apache.cassandra.schema.LegacySchemaTables; import org.apache.cassandra.db.SystemKeyspace; import org.apache.cassandra.exceptions.AuthenticationException; import org.apache.cassandra.exceptions.InvalidRequestException; import org.apache.cassandra.exceptions.UnauthorizedException; -import org.apache.cassandra.tracing.TraceKeyspace; +import org.apache.cassandra.schema.LegacySchemaTables; import org.apache.cassandra.thrift.ThriftValidation; +import org.apache.cassandra.tracing.TraceKeyspace; import org.apache.cassandra.utils.FBUtilities; import org.apache.cassandra.utils.JVMStabilityInspector; import org.apache.cassandra.utils.SemanticVersion; @@ -52,16 +53,27 @@ public class ClientState private static final Set READABLE_SYSTEM_RESOURCES = new HashSet<>(); private static final Set PROTECTED_AUTH_RESOURCES = new HashSet<>(); - + private static final Set ALTERABLE_SYSTEM_KEYSPACES = new HashSet<>(); + private static final Set DROPPABLE_SYSTEM_TABLES = new HashSet<>(); static { // We want these system cfs to be always readable to authenticated users since many tools rely on them // (nodetool, cqlsh, bulkloader, etc.) for (String cf : Iterables.concat(Arrays.asList(SystemKeyspace.LOCAL, SystemKeyspace.PEERS), LegacySchemaTables.ALL)) - READABLE_SYSTEM_RESOURCES.add(DataResource.columnFamily(SystemKeyspace.NAME, cf)); + READABLE_SYSTEM_RESOURCES.add(DataResource.table(SystemKeyspace.NAME, cf)); PROTECTED_AUTH_RESOURCES.addAll(DatabaseDescriptor.getAuthenticator().protectedResources()); PROTECTED_AUTH_RESOURCES.addAll(DatabaseDescriptor.getAuthorizer().protectedResources()); + PROTECTED_AUTH_RESOURCES.addAll(DatabaseDescriptor.getRoleManager().protectedResources()); + + // allow users with sufficient privileges to alter KS level options on AUTH_KS and + // TRACING_KS, and also to drop legacy tables (users, credentials, permissions) from + // AUTH_KS + ALTERABLE_SYSTEM_KEYSPACES.add(AuthKeyspace.NAME); + ALTERABLE_SYSTEM_KEYSPACES.add(TraceKeyspace.NAME); + DROPPABLE_SYSTEM_TABLES.add(DataResource.table(AuthKeyspace.NAME, PasswordAuthenticator.LEGACY_CREDENTIALS_TABLE)); + DROPPABLE_SYSTEM_TABLES.add(DataResource.table(AuthKeyspace.NAME, CassandraRoleManager.LEGACY_USERS_TABLE)); + DROPPABLE_SYSTEM_TABLES.add(DataResource.table(AuthKeyspace.NAME, CassandraAuthorizer.USER_PERMISSIONS)); } // Current user for the session @@ -200,10 +212,13 @@ public class ClientState */ public void login(AuthenticatedUser user) throws AuthenticationException { - if (!user.isAnonymous() && !Auth.isExistingUser(user.getName())) - throw new AuthenticationException(String.format("User %s doesn't exist - create it with CREATE USER query first", - user.getName())); - this.user = user; + // Login privilege is not inherited via granted roles, so just + // verify that the role with the credentials that were actually + // supplied has it + if (user.isAnonymous() || DatabaseDescriptor.getRoleManager().canLogin(user.getName())) + this.user = user; + else + throw new AuthenticationException(String.format("%s is not permitted to log in", user.getName())); } public void hasAllKeyspacesAccess(Permission perm) throws UnauthorizedException @@ -223,7 +238,7 @@ public class ClientState throws UnauthorizedException, InvalidRequestException { ThriftValidation.validateColumnFamily(keyspace, columnFamily); - hasAccess(keyspace, perm, DataResource.columnFamily(keyspace, columnFamily)); + hasAccess(keyspace, perm, DataResource.table(keyspace, columnFamily)); } private void hasAccess(String keyspace, Permission perm, DataResource resource) @@ -264,10 +279,15 @@ public class ClientState if (SystemKeyspace.NAME.equalsIgnoreCase(keyspace)) throw new UnauthorizedException(keyspace + " keyspace is not user-modifiable."); - // we want to allow altering AUTH_KS and TRACING_KS. - Set allowAlter = Sets.newHashSet(Auth.AUTH_KS, TraceKeyspace.NAME); - if (allowAlter.contains(keyspace.toLowerCase()) && !(resource.isKeyspaceLevel() && (perm == Permission.ALTER))) + // allow users with sufficient privileges to alter KS level options on AUTH_KS and + // TRACING_KS, and also to drop legacy tables (users, credentials, permissions) from + // AUTH_KS + if (ALTERABLE_SYSTEM_KEYSPACES.contains(resource.getKeyspace().toLowerCase()) + && ((perm == Permission.ALTER && !resource.isKeyspaceLevel()) + || (perm == Permission.DROP && !DROPPABLE_SYSTEM_TABLES.contains(resource)))) + { throw new UnauthorizedException(String.format("Cannot %s %s", perm, resource)); + } } public void validateLogin() throws UnauthorizedException @@ -307,6 +327,6 @@ public class ClientState private Set authorize(IResource resource) { - return Auth.getPermissions(user, resource); + return AuthenticatedUser.getPermissions(user, resource); } } diff --git a/src/java/org/apache/cassandra/service/StorageService.java b/src/java/org/apache/cassandra/service/StorageService.java index 4740cd3ffb..bb3e88294c 100644 --- a/src/java/org/apache/cassandra/service/StorageService.java +++ b/src/java/org/apache/cassandra/service/StorageService.java @@ -17,10 +17,7 @@ */ package org.apache.cassandra.service; -import java.io.ByteArrayInputStream; -import java.io.DataInputStream; -import java.io.File; -import java.io.IOException; +import java.io.*; import java.lang.management.ManagementFactory; import java.net.InetAddress; import java.net.UnknownHostException; @@ -29,20 +26,10 @@ import java.util.*; import java.util.concurrent.*; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicLong; - -import javax.management.JMX; -import javax.management.MBeanServer; -import javax.management.Notification; -import javax.management.NotificationBroadcasterSupport; -import javax.management.ObjectName; +import javax.management.*; import javax.management.openmbean.TabularData; import javax.management.openmbean.TabularDataSupport; -import ch.qos.logback.classic.LoggerContext; -import ch.qos.logback.classic.jmx.JMXConfiguratorMBean; -import ch.qos.logback.classic.spi.ILoggingEvent; -import ch.qos.logback.core.Appender; - import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Predicate; import com.google.common.collect.*; @@ -52,12 +39,17 @@ import org.apache.commons.lang3.time.DurationFormatUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.apache.cassandra.auth.Auth; +import ch.qos.logback.classic.LoggerContext; +import ch.qos.logback.classic.jmx.JMXConfiguratorMBean; +import ch.qos.logback.classic.spi.ILoggingEvent; +import ch.qos.logback.core.Appender; +import org.apache.cassandra.auth.AuthKeyspace; +import org.apache.cassandra.auth.AuthMigrationListener; import org.apache.cassandra.concurrent.*; import org.apache.cassandra.config.*; -import org.apache.cassandra.cql3.UntypedResultSet; import org.apache.cassandra.cql3.QueryOptions; import org.apache.cassandra.cql3.QueryProcessor; +import org.apache.cassandra.cql3.UntypedResultSet; import org.apache.cassandra.cql3.statements.SelectStatement; import org.apache.cassandra.db.*; import org.apache.cassandra.db.commitlog.CommitLog; @@ -74,15 +66,9 @@ import org.apache.cassandra.io.sstable.SSTableLoader; import org.apache.cassandra.io.util.FileUtils; import org.apache.cassandra.locator.*; import org.apache.cassandra.metrics.StorageMetrics; -import org.apache.cassandra.net.AsyncOneResponse; -import org.apache.cassandra.net.MessageOut; -import org.apache.cassandra.net.MessagingService; -import org.apache.cassandra.net.ResponseVerbHandler; -import org.apache.cassandra.repair.RepairMessageVerbHandler; -import org.apache.cassandra.repair.RepairSessionResult; +import org.apache.cassandra.net.*; +import org.apache.cassandra.repair.*; import org.apache.cassandra.repair.messages.RepairOption; -import org.apache.cassandra.repair.RepairSession; -import org.apache.cassandra.repair.RepairParallelism; import org.apache.cassandra.service.paxos.CommitVerbHandler; import org.apache.cassandra.service.paxos.PrepareVerbHandler; import org.apache.cassandra.service.paxos.ProposeVerbHandler; @@ -843,7 +829,7 @@ public class StorageService extends NotificationBroadcasterSupport implements IE Gossiper.instance.replacedEndpoint(existing); assert tokenMetadata.sortedTokens().size() > 0; - Auth.setup(); + doAuthSetup(); } else { @@ -882,10 +868,41 @@ public class StorageService extends NotificationBroadcasterSupport implements IE logger.info("Leaving write survey mode and joining ring at operator request"); assert tokenMetadata.sortedTokens().size() > 0; - Auth.setup(); + doAuthSetup(); } } + private void doAuthSetup() + { + try + { + // if we don't have system_auth keyspace at this point, then create it manually + // otherwise, create any necessary tables as we may be upgrading in which case + // the ks exists with the only the legacy tables defined + if (Schema.instance.getKSMetaData(AuthKeyspace.NAME) == null) + { + MigrationManager.announceNewKeyspace(AuthKeyspace.definition(), 0, false); + } + else + { + for (Map.Entry table : AuthKeyspace.definition().cfMetaData().entrySet()) + { + if (Schema.instance.getCFMetaData(AuthKeyspace.NAME, table.getKey()) == null) + MigrationManager.announceNewColumnFamily(table.getValue()); + } + } + } + catch (Exception e) + { + throw new AssertionError(e); // shouldn't ever happen. + } + + DatabaseDescriptor.getRoleManager().setup(); + DatabaseDescriptor.getAuthenticator().setup(); + DatabaseDescriptor.getAuthorizer().setup(); + MigrationManager.instance.register(new AuthMigrationListener()); + } + public boolean isJoined() { return joined; diff --git a/src/java/org/apache/cassandra/thrift/CassandraServer.java b/src/java/org/apache/cassandra/thrift/CassandraServer.java index 7d89049348..c0de59f311 100644 --- a/src/java/org/apache/cassandra/thrift/CassandraServer.java +++ b/src/java/org/apache/cassandra/thrift/CassandraServer.java @@ -29,43 +29,30 @@ import java.util.zip.Inflater; import com.google.common.base.Function; import com.google.common.base.Joiner; -import com.google.common.collect.ImmutableMap; -import com.google.common.collect.ImmutableSortedSet; -import com.google.common.collect.Iterables; -import com.google.common.collect.Lists; -import com.google.common.collect.Maps; +import com.google.common.collect.*; import com.google.common.primitives.Longs; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.apache.cassandra.auth.AuthenticatedUser; import org.apache.cassandra.auth.Permission; -import org.apache.cassandra.config.CFMetaData; -import org.apache.cassandra.config.DatabaseDescriptor; -import org.apache.cassandra.config.KSMetaData; -import org.apache.cassandra.config.Schema; +import org.apache.cassandra.config.*; import org.apache.cassandra.cql3.QueryOptions; import org.apache.cassandra.cql3.statements.ParsedStatement; import org.apache.cassandra.db.*; import org.apache.cassandra.db.composites.*; import org.apache.cassandra.db.context.CounterContext; import org.apache.cassandra.db.filter.ColumnSlice; -import org.apache.cassandra.db.filter.IDiskAtomFilter; -import org.apache.cassandra.db.filter.NamesQueryFilter; -import org.apache.cassandra.db.filter.SliceQueryFilter; +import org.apache.cassandra.db.filter.*; import org.apache.cassandra.db.marshal.TimeUUIDType; import org.apache.cassandra.dht.*; +import org.apache.cassandra.dht.Range; import org.apache.cassandra.exceptions.*; import org.apache.cassandra.io.util.DataOutputBuffer; import org.apache.cassandra.locator.DynamicEndpointSnitch; import org.apache.cassandra.metrics.ClientMetrics; import org.apache.cassandra.scheduler.IRequestScheduler; import org.apache.cassandra.serializers.MarshalException; -import org.apache.cassandra.service.CASRequest; -import org.apache.cassandra.service.ClientState; -import org.apache.cassandra.service.MigrationManager; -import org.apache.cassandra.service.StorageProxy; -import org.apache.cassandra.service.StorageService; +import org.apache.cassandra.service.*; import org.apache.cassandra.service.pager.QueryPagers; import org.apache.cassandra.tracing.Tracing; import org.apache.cassandra.utils.ByteBufferUtil; @@ -1472,12 +1459,11 @@ public class CassandraServer implements Cassandra.Iface } } - public void login(AuthenticationRequest auth_request) throws AuthenticationException, AuthorizationException, TException + public void login(AuthenticationRequest auth_request) throws TException { try { - AuthenticatedUser user = DatabaseDescriptor.getAuthenticator().authenticate(auth_request.getCredentials()); - state().login(user); + state().login(DatabaseDescriptor.getAuthenticator().legacyAuthenticate(auth_request.getCredentials())); } catch (org.apache.cassandra.exceptions.AuthenticationException e) { diff --git a/src/java/org/apache/cassandra/tools/BulkLoader.java b/src/java/org/apache/cassandra/tools/BulkLoader.java index a720e12264..f89acfb5d0 100644 --- a/src/java/org/apache/cassandra/tools/BulkLoader.java +++ b/src/java/org/apache/cassandra/tools/BulkLoader.java @@ -18,33 +18,33 @@ package org.apache.cassandra.tools; import java.io.File; -import java.net.*; +import java.net.InetAddress; +import java.net.MalformedURLException; +import java.net.UnknownHostException; import java.util.*; import com.google.common.base.Joiner; import com.google.common.collect.HashMultimap; import com.google.common.collect.Multimap; - import org.apache.commons.cli.*; -import org.apache.thrift.protocol.TBinaryProtocol; -import org.apache.thrift.protocol.TProtocol; -import org.apache.thrift.transport.TTransport; - -import org.apache.cassandra.auth.IAuthenticator; +import org.apache.cassandra.auth.PasswordAuthenticator; import org.apache.cassandra.config.*; -import org.apache.cassandra.schema.LegacySchemaTables; import org.apache.cassandra.db.SystemKeyspace; import org.apache.cassandra.db.marshal.UTF8Type; import org.apache.cassandra.dht.Range; import org.apache.cassandra.dht.Token; import org.apache.cassandra.exceptions.ConfigurationException; import org.apache.cassandra.io.sstable.SSTableLoader; +import org.apache.cassandra.schema.LegacySchemaTables; import org.apache.cassandra.streaming.*; import org.apache.cassandra.thrift.*; import org.apache.cassandra.utils.ByteBufferUtil; import org.apache.cassandra.utils.JVMStabilityInspector; import org.apache.cassandra.utils.OutputHandler; +import org.apache.thrift.protocol.TBinaryProtocol; +import org.apache.thrift.protocol.TProtocol; +import org.apache.thrift.transport.TTransport; public class BulkLoader { @@ -359,8 +359,8 @@ public class BulkLoader if (user != null && passwd != null) { Map credentials = new HashMap<>(); - credentials.put(IAuthenticator.USERNAME_KEY, user); - credentials.put(IAuthenticator.PASSWORD_KEY, passwd); + credentials.put(PasswordAuthenticator.USERNAME_KEY, user); + credentials.put(PasswordAuthenticator.PASSWORD_KEY, passwd); AuthenticationRequest authenticationRequest = new AuthenticationRequest(credentials); client.login(authenticationRequest); } diff --git a/src/java/org/apache/cassandra/transport/Client.java b/src/java/org/apache/cassandra/transport/Client.java index c7c3103953..571a7ce8bc 100644 --- a/src/java/org/apache/cassandra/transport/Client.java +++ b/src/java/org/apache/cassandra/transport/Client.java @@ -22,16 +22,11 @@ import java.io.IOException; import java.io.InputStreamReader; import java.nio.ByteBuffer; import java.nio.charset.StandardCharsets; -import java.util.ArrayList; -import java.util.Collections; -import java.util.HashMap; -import java.util.Iterator; -import java.util.List; -import java.util.Map; +import java.util.*; import com.google.common.base.Splitter; -import org.apache.cassandra.auth.IAuthenticator; +import org.apache.cassandra.auth.PasswordAuthenticator; import org.apache.cassandra.cql3.QueryOptions; import org.apache.cassandra.db.ConsistencyLevel; import org.apache.cassandra.db.marshal.Int32Type; @@ -179,7 +174,7 @@ public class Client extends SimpleClient else if (msgType.equals("AUTHENTICATE")) { Map credentials = readCredentials(iter); - if(!credentials.containsKey(IAuthenticator.USERNAME_KEY) || !credentials.containsKey(IAuthenticator.PASSWORD_KEY)) + if(!credentials.containsKey(PasswordAuthenticator.USERNAME_KEY) || !credentials.containsKey(PasswordAuthenticator.PASSWORD_KEY)) { System.err.println("[ERROR] Authentication requires both 'username' and 'password'"); return null; @@ -221,8 +216,8 @@ public class Client extends SimpleClient private byte[] encodeCredentialsForSasl(Map credentials) { - byte[] username = credentials.get(IAuthenticator.USERNAME_KEY).getBytes(StandardCharsets.UTF_8); - byte[] password = credentials.get(IAuthenticator.PASSWORD_KEY).getBytes(StandardCharsets.UTF_8); + byte[] username = credentials.get(PasswordAuthenticator.USERNAME_KEY).getBytes(StandardCharsets.UTF_8); + byte[] password = credentials.get(PasswordAuthenticator.PASSWORD_KEY).getBytes(StandardCharsets.UTF_8); byte[] initialResponse = new byte[username.length + password.length + 2]; initialResponse[0] = 0; System.arraycopy(username, 0, initialResponse, 1, username.length); diff --git a/src/java/org/apache/cassandra/transport/Server.java b/src/java/org/apache/cassandra/transport/Server.java index 99601a6ce3..8830479222 100644 --- a/src/java/org/apache/cassandra/transport/Server.java +++ b/src/java/org/apache/cassandra/transport/Server.java @@ -28,21 +28,24 @@ import java.util.concurrent.atomic.AtomicBoolean; import javax.net.ssl.SSLContext; import javax.net.ssl.SSLEngine; -import io.netty.channel.epoll.Epoll; -import io.netty.channel.epoll.EpollEventLoopGroup; -import io.netty.channel.epoll.EpollServerSocketChannel; -import io.netty.util.Version; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import io.netty.bootstrap.ServerBootstrap; +import io.netty.channel.*; +import io.netty.channel.epoll.Epoll; +import io.netty.channel.epoll.EpollEventLoopGroup; +import io.netty.channel.epoll.EpollServerSocketChannel; +import io.netty.channel.group.ChannelGroup; +import io.netty.channel.group.DefaultChannelGroup; import io.netty.channel.nio.NioEventLoopGroup; import io.netty.channel.socket.nio.NioServerSocketChannel; +import io.netty.handler.ssl.SslHandler; +import io.netty.util.Version; import io.netty.util.concurrent.EventExecutor; import io.netty.util.concurrent.GlobalEventExecutor; import io.netty.util.internal.logging.InternalLoggerFactory; import io.netty.util.internal.logging.Slf4JLoggerFactory; -import org.apache.cassandra.auth.IAuthenticator; -import org.apache.cassandra.auth.ISaslAwareAuthenticator; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.config.EncryptionOptions; import org.apache.cassandra.db.marshal.AbstractType; @@ -50,11 +53,6 @@ import org.apache.cassandra.metrics.ClientMetrics; import org.apache.cassandra.security.SSLFactory; import org.apache.cassandra.service.*; import org.apache.cassandra.transport.messages.EventMessage; -import io.netty.bootstrap.ServerBootstrap; -import io.netty.channel.*; -import io.netty.channel.group.ChannelGroup; -import io.netty.channel.group.DefaultChannelGroup; -import io.netty.handler.ssl.SslHandler; public class Server implements CassandraDaemon.Server { @@ -132,16 +130,6 @@ public class Server implements CassandraDaemon.Server private void run() { - // Check that a SaslAuthenticator can be provided by the configured - // IAuthenticator. If not, don't start the server. - IAuthenticator authenticator = DatabaseDescriptor.getAuthenticator(); - if (authenticator.requireAuthentication() && !(authenticator instanceof ISaslAwareAuthenticator)) - { - logger.error("Not starting native transport as the configured IAuthenticator is not capable of SASL authentication"); - isRunning.compareAndSet(true, false); - return; - } - // Configure the server. eventExecutorGroup = new RequestThreadPoolExecutor(); diff --git a/src/java/org/apache/cassandra/transport/ServerConnection.java b/src/java/org/apache/cassandra/transport/ServerConnection.java index b28866f85a..24eb6437a3 100644 --- a/src/java/org/apache/cassandra/transport/ServerConnection.java +++ b/src/java/org/apache/cassandra/transport/ServerConnection.java @@ -20,21 +20,17 @@ package org.apache.cassandra.transport; import java.util.concurrent.ConcurrentMap; import io.netty.channel.Channel; - import org.apache.cassandra.auth.IAuthenticator; -import org.apache.cassandra.auth.ISaslAwareAuthenticator; -import org.apache.cassandra.auth.ISaslAwareAuthenticator.SaslAuthenticator; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.service.ClientState; import org.apache.cassandra.service.QueryState; - import org.cliffc.high_scale_lib.NonBlockingHashMap; public class ServerConnection extends Connection { private enum State { UNINITIALIZED, AUTHENTICATION, READY } - private volatile SaslAuthenticator saslAuthenticator; + private volatile IAuthenticator.SaslNegotiator saslNegotiator; private final ClientState clientState; private volatile State state; @@ -104,7 +100,7 @@ public class ServerConnection extends Connection { state = State.READY; // we won't use the authenticator again, null it so that it can be GC'd - saslAuthenticator = null; + saslNegotiator = null; } break; case READY: @@ -114,14 +110,10 @@ public class ServerConnection extends Connection } } - public SaslAuthenticator getAuthenticator() + public IAuthenticator.SaslNegotiator getSaslNegotiator() { - if (saslAuthenticator == null) - { - IAuthenticator authenticator = DatabaseDescriptor.getAuthenticator(); - assert authenticator instanceof ISaslAwareAuthenticator : "Configured IAuthenticator does not support SASL authentication"; - saslAuthenticator = ((ISaslAwareAuthenticator)authenticator).newAuthenticator(); - } - return saslAuthenticator; + if (saslNegotiator == null) + saslNegotiator = DatabaseDescriptor.getAuthenticator().newSaslNegotiator(); + return saslNegotiator; } } diff --git a/src/java/org/apache/cassandra/transport/messages/AuthResponse.java b/src/java/org/apache/cassandra/transport/messages/AuthResponse.java index 3f3f7745a0..cb67476dd8 100644 --- a/src/java/org/apache/cassandra/transport/messages/AuthResponse.java +++ b/src/java/org/apache/cassandra/transport/messages/AuthResponse.java @@ -17,18 +17,14 @@ */ package org.apache.cassandra.transport.messages; -import org.apache.cassandra.auth.AuthenticatedUser; -import org.apache.cassandra.auth.ISaslAwareAuthenticator.SaslAuthenticator; -import org.apache.cassandra.exceptions.AuthenticationException; -import org.apache.cassandra.service.QueryState; -import org.apache.cassandra.transport.CBUtil; -import org.apache.cassandra.transport.Message; -import org.apache.cassandra.transport.ProtocolException; -import org.apache.cassandra.transport.ServerConnection; +import java.nio.ByteBuffer; import io.netty.buffer.ByteBuf; - -import java.nio.ByteBuffer; +import org.apache.cassandra.auth.AuthenticatedUser; +import org.apache.cassandra.auth.IAuthenticator; +import org.apache.cassandra.exceptions.AuthenticationException; +import org.apache.cassandra.service.QueryState; +import org.apache.cassandra.transport.*; /** * A SASL token message sent from client to server. Some SASL @@ -61,11 +57,12 @@ public class AuthResponse extends Message.Request } }; - private byte[] token; + private final byte[] token; public AuthResponse(byte[] token) { super(Message.Type.AUTH_RESPONSE); + assert token != null; this.token = token; } @@ -74,11 +71,11 @@ public class AuthResponse extends Message.Request { try { - SaslAuthenticator authenticator = ((ServerConnection) connection).getAuthenticator(); - byte[] challenge = authenticator.evaluateResponse(token == null ? new byte[0] : token); - if (authenticator.isComplete()) + IAuthenticator.SaslNegotiator negotiator = ((ServerConnection) connection).getSaslNegotiator(); + byte[] challenge = negotiator.evaluateResponse(token); + if (negotiator.isComplete()) { - AuthenticatedUser user = authenticator.getAuthenticatedUser(); + AuthenticatedUser user = negotiator.getAuthenticatedUser(); queryState.getClientState().login(user); // authentication is complete, send a ready message to the client return new AuthSuccess(challenge); diff --git a/src/java/org/apache/cassandra/transport/messages/CredentialsMessage.java b/src/java/org/apache/cassandra/transport/messages/CredentialsMessage.java index eb39e30b75..aad4232f65 100644 --- a/src/java/org/apache/cassandra/transport/messages/CredentialsMessage.java +++ b/src/java/org/apache/cassandra/transport/messages/CredentialsMessage.java @@ -20,15 +20,14 @@ package org.apache.cassandra.transport.messages; import java.util.HashMap; import java.util.Map; +import io.netty.buffer.ByteBuf; import org.apache.cassandra.auth.AuthenticatedUser; import org.apache.cassandra.config.DatabaseDescriptor; -import org.apache.cassandra.transport.ProtocolException; -import io.netty.buffer.ByteBuf; - import org.apache.cassandra.exceptions.AuthenticationException; import org.apache.cassandra.service.QueryState; import org.apache.cassandra.transport.CBUtil; import org.apache.cassandra.transport.Message; +import org.apache.cassandra.transport.ProtocolException; /** * Message to indicate that the server is ready to receive requests. @@ -75,14 +74,15 @@ public class CredentialsMessage extends Message.Request { try { - AuthenticatedUser user = DatabaseDescriptor.getAuthenticator().authenticate(credentials); + AuthenticatedUser user = DatabaseDescriptor.getAuthenticator().legacyAuthenticate(credentials); state.getClientState().login(user); - return new ReadyMessage(); } catch (AuthenticationException e) { return ErrorMessage.fromException(e); } + + return new ReadyMessage(); } @Override diff --git a/src/java/org/apache/cassandra/utils/FBUtilities.java b/src/java/org/apache/cassandra/utils/FBUtilities.java index 8077df8e5e..0462e5e68b 100644 --- a/src/java/org/apache/cassandra/utils/FBUtilities.java +++ b/src/java/org/apache/cassandra/utils/FBUtilities.java @@ -20,19 +20,12 @@ package org.apache.cassandra.utils; import java.io.*; import java.lang.reflect.Field; import java.math.BigInteger; -import java.net.InetAddress; -import java.net.NetworkInterface; -import java.net.SocketException; -import java.net.URL; -import java.net.UnknownHostException; +import java.net.*; import java.nio.ByteBuffer; import java.security.MessageDigest; import java.security.NoSuchAlgorithmException; import java.util.*; -import java.util.concurrent.ExecutionException; -import java.util.concurrent.Future; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.TimeoutException; +import java.util.concurrent.*; import java.util.zip.Checksum; import com.google.common.base.Joiner; @@ -43,6 +36,7 @@ import org.slf4j.LoggerFactory; import org.apache.cassandra.auth.IAuthenticator; import org.apache.cassandra.auth.IAuthorizer; +import org.apache.cassandra.auth.IRoleManager; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.DecoratedKey; import org.apache.cassandra.dht.IPartitioner; @@ -54,10 +48,7 @@ import org.apache.cassandra.io.util.DataOutputBuffer; import org.apache.cassandra.io.util.FileUtils; import org.apache.cassandra.io.util.IAllocator; import org.apache.cassandra.net.AsyncOneResponse; -import org.apache.thrift.TBase; -import org.apache.thrift.TDeserializer; -import org.apache.thrift.TException; -import org.apache.thrift.TSerializer; +import org.apache.thrift.*; import org.codehaus.jackson.JsonFactory; import org.codehaus.jackson.map.ObjectMapper; @@ -439,6 +430,13 @@ public class FBUtilities return FBUtilities.construct(className, "authenticator"); } + public static IRoleManager newRoleManager(String className) throws ConfigurationException + { + if (!className.contains(".")) + className = "org.apache.cassandra.auth." + className; + return FBUtilities.construct(className, "role manager"); + } + /** * @return The Class for the given name. * @param classname Fully qualified classname.