Updated Branches: refs/heads/trunk 822ee88a3 -> c74a871e7
Fix bug starting C* with simple authentication patch by slebresne; reviewed by jbellis for CASSANDRA-4648 Project: http://git-wip-us.apache.org/repos/asf/cassandra/repo Commit: http://git-wip-us.apache.org/repos/asf/cassandra/commit/c74a871e Tree: http://git-wip-us.apache.org/repos/asf/cassandra/tree/c74a871e Diff: http://git-wip-us.apache.org/repos/asf/cassandra/diff/c74a871e Branch: refs/heads/trunk Commit: c74a871e7a5882802de46550ea5bab0220f16883 Parents: 822ee88 Author: Sylvain Lebresne <[email protected]> Authored: Wed Sep 12 14:21:17 2012 +0200 Committer: Sylvain Lebresne <[email protected]> Committed: Thu Sep 27 11:02:03 2012 +0200 ---------------------------------------------------------------------- CHANGES.txt | 2 + .../org/apache/cassandra/cql3/CQLStatement.java | 7 + .../org/apache/cassandra/cql3/QueryProcessor.java | 5 +- .../cassandra/cql3/statements/BatchStatement.java | 4 +- .../cassandra/cql3/statements/CFStatement.java | 1 + .../cassandra/cql3/statements/DeleteStatement.java | 4 +- .../cassandra/cql3/statements/GrantStatement.java | 19 +--- .../cql3/statements/ListGrantsStatement.java | 15 +-- .../cql3/statements/ModificationStatement.java | 19 ++- .../statements/PermissionAlteringStatement.java | 48 ++++++ .../cassandra/cql3/statements/RevokeStatement.java | 21 +--- .../cql3/statements/SchemaAlteringStatement.java | 11 +- .../cassandra/cql3/statements/SelectStatement.java | 111 +++++++++------ .../cql3/statements/TruncateStatement.java | 5 + .../cassandra/cql3/statements/UpdateStatement.java | 4 +- .../cassandra/cql3/statements/UseStatement.java | 6 + .../cassandra/service/RangeSliceVerbHandler.java | 2 +- 17 files changed, 170 insertions(+), 114 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/cassandra/blob/c74a871e/CHANGES.txt ---------------------------------------------------------------------- diff --git a/CHANGES.txt b/CHANGES.txt index 4d8eb91..576f7a6 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -7,6 +7,8 @@ * Fix scrubbing of CQL3 created tables (CASSANDRA-4685) * (CQL3) Fix validation when using counter and regular columns in the same table (CASSANDRA-4706) + * Fix bug starting Cassandra with simple authentication (CASSANDRA-4648 + * ) 1.2-beta1 * add atomic_batch_mutate (CASSANDRA-4542, -4635) http://git-wip-us.apache.org/repos/asf/cassandra/blob/c74a871e/src/java/org/apache/cassandra/cql3/CQLStatement.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/cql3/CQLStatement.java b/src/java/org/apache/cassandra/cql3/CQLStatement.java index d0889bc..00fb406 100644 --- a/src/java/org/apache/cassandra/cql3/CQLStatement.java +++ b/src/java/org/apache/cassandra/cql3/CQLStatement.java @@ -54,4 +54,11 @@ public interface CQLStatement * can assume that each bound term have a corresponding value. */ public ResultMessage execute(ClientState state, List<ByteBuffer> variables) throws RequestValidationException, RequestExecutionException; + + /** + * Variante of execute used for internal query against the system tables, and thus only query the local node. + * + * @param state the current client state + */ + public ResultMessage executeInternal(ClientState state) throws RequestValidationException, RequestExecutionException; } http://git-wip-us.apache.org/repos/asf/cassandra/blob/c74a871e/src/java/org/apache/cassandra/cql3/QueryProcessor.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/cql3/QueryProcessor.java b/src/java/org/apache/cassandra/cql3/QueryProcessor.java index 856f6fd..58a8f44 100644 --- a/src/java/org/apache/cassandra/cql3/QueryProcessor.java +++ b/src/java/org/apache/cassandra/cql3/QueryProcessor.java @@ -143,7 +143,10 @@ public class QueryProcessor try { ClientState state = new ClientState(true); - ResultMessage result = processStatement(getStatement(query, state).statement, state, Collections.<ByteBuffer>emptyList()); + state.setKeyspace(Table.SYSTEM_KS); + CQLStatement statement = getStatement(query, state).statement; + statement.validate(state); + ResultMessage result = statement.executeInternal(state); if (result instanceof ResultMessage.Rows) return new UntypedResultSet(((ResultMessage.Rows)result).result); else http://git-wip-us.apache.org/repos/asf/cassandra/blob/c74a871e/src/java/org/apache/cassandra/cql3/statements/BatchStatement.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/cql3/statements/BatchStatement.java b/src/java/org/apache/cassandra/cql3/statements/BatchStatement.java index 246a97b..9e5cb40 100644 --- a/src/java/org/apache/cassandra/cql3/statements/BatchStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/BatchStatement.java @@ -119,7 +119,7 @@ public class BatchStatement extends ModificationStatement } } - public List<IMutation> getMutations(ClientState clientState, List<ByteBuffer> variables) + public List<IMutation> getMutations(ClientState clientState, List<ByteBuffer> variables, boolean local) throws RequestExecutionException, RequestValidationException { Map<Pair<String, ByteBuffer>, RowAndCounterMutation> mutations = new HashMap<Pair<String, ByteBuffer>, RowAndCounterMutation>(); @@ -128,7 +128,7 @@ public class BatchStatement extends ModificationStatement if (isSetTimestamp()) statement.setTimestamp(getTimestamp(clientState)); - List<IMutation> lm = statement.getMutations(clientState, variables); + List<IMutation> lm = statement.getMutations(clientState, variables, local); // Group mutation together, otherwise they won't get applied atomically for (IMutation m : lm) { http://git-wip-us.apache.org/repos/asf/cassandra/blob/c74a871e/src/java/org/apache/cassandra/cql3/statements/CFStatement.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/cql3/statements/CFStatement.java b/src/java/org/apache/cassandra/cql3/statements/CFStatement.java index e9244ab..2ccc203 100644 --- a/src/java/org/apache/cassandra/cql3/statements/CFStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/CFStatement.java @@ -37,6 +37,7 @@ public abstract class CFStatement extends ParsedStatement { if (!cfName.hasKeyspace()) { + // XXX: We explicitely only want to call state.getKeyspace() in this case, don't move it outside the if. cfName.setKeyspace(state.getKeyspace(), true); } } http://git-wip-us.apache.org/repos/asf/cassandra/blob/c74a871e/src/java/org/apache/cassandra/cql3/statements/DeleteStatement.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/cql3/statements/DeleteStatement.java b/src/java/org/apache/cassandra/cql3/statements/DeleteStatement.java index 73f7c0f..50a279f 100644 --- a/src/java/org/apache/cassandra/cql3/statements/DeleteStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/DeleteStatement.java @@ -62,7 +62,7 @@ public class DeleteStatement extends ModificationStatement this.toRemove = new ArrayList<Pair<CFDefinition.Name, Term>>(columns.size()); } - public List<IMutation> getMutations(ClientState clientState, List<ByteBuffer> variables) + public List<IMutation> getMutations(ClientState clientState, List<ByteBuffer> variables, boolean local) throws RequestExecutionException, RequestValidationException { // keys @@ -92,7 +92,7 @@ public class DeleteStatement extends ModificationStatement } } - Map<ByteBuffer, ColumnGroupMap> rows = needsReading ? readRows(keys, builder, (CompositeType)cfDef.cfm.comparator) : null; + Map<ByteBuffer, ColumnGroupMap> rows = needsReading ? readRows(keys, builder, (CompositeType)cfDef.cfm.comparator, local) : null; List<IMutation> rowMutations = new ArrayList<IMutation>(keys.size()); UpdateParameters params = new UpdateParameters(variables, getTimestamp(clientState), -1); http://git-wip-us.apache.org/repos/asf/cassandra/blob/c74a871e/src/java/org/apache/cassandra/cql3/statements/GrantStatement.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/cql3/statements/GrantStatement.java b/src/java/org/apache/cassandra/cql3/statements/GrantStatement.java index 365b9cd..45d7d2a 100644 --- a/src/java/org/apache/cassandra/cql3/statements/GrantStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/GrantStatement.java @@ -23,13 +23,12 @@ import java.util.List; import org.apache.cassandra.auth.Permission; import org.apache.cassandra.cql3.CFName; -import org.apache.cassandra.cql3.CQLStatement; import org.apache.cassandra.exceptions.InvalidRequestException; import org.apache.cassandra.exceptions.UnauthorizedException; import org.apache.cassandra.service.ClientState; import org.apache.cassandra.transport.messages.ResultMessage; -public class GrantStatement extends ParsedStatement implements CQLStatement +public class GrantStatement extends PermissionAlteringStatement { private final Permission permission; private final CFName resource; @@ -44,25 +43,9 @@ public class GrantStatement extends ParsedStatement implements CQLStatement this.grantOption = grantOption; } - public int getBoundsTerms() - { - return 0; - } - - public void checkAccess(ClientState state) throws UnauthorizedException, InvalidRequestException - {} - - public void validate(ClientState state) throws InvalidRequestException - {} - public ResultMessage execute(ClientState state, List<ByteBuffer> variables) throws UnauthorizedException, InvalidRequestException { state.grantPermission(permission, username, resource, grantOption); return null; } - - public Prepared prepare() throws InvalidRequestException - { - return new Prepared(this); - } } http://git-wip-us.apache.org/repos/asf/cassandra/blob/c74a871e/src/java/org/apache/cassandra/cql3/statements/ListGrantsStatement.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/cql3/statements/ListGrantsStatement.java b/src/java/org/apache/cassandra/cql3/statements/ListGrantsStatement.java index e929d36..ecc312f 100644 --- a/src/java/org/apache/cassandra/cql3/statements/ListGrantsStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/ListGrantsStatement.java @@ -21,13 +21,12 @@ package org.apache.cassandra.cql3.statements; import java.nio.ByteBuffer; import java.util.List; -import org.apache.cassandra.cql3.CQLStatement; import org.apache.cassandra.exceptions.InvalidRequestException; import org.apache.cassandra.exceptions.UnauthorizedException; import org.apache.cassandra.service.ClientState; import org.apache.cassandra.transport.messages.ResultMessage; -public class ListGrantsStatement extends ParsedStatement implements CQLStatement +public class ListGrantsStatement extends PermissionAlteringStatement { private final String username; @@ -36,20 +35,8 @@ public class ListGrantsStatement extends ParsedStatement implements CQLStatement this.username = username; } - public void checkAccess(ClientState state) throws InvalidRequestException - {} - - public void validate(ClientState state) throws InvalidRequestException - {} - public ResultMessage execute(ClientState state, List<ByteBuffer> variables) throws UnauthorizedException, InvalidRequestException { return state.listPermissions(username); } - - @Override - public Prepared prepare() throws InvalidRequestException - { - return new Prepared(this); - } } http://git-wip-us.apache.org/repos/asf/cassandra/blob/c74a871e/src/java/org/apache/cassandra/cql3/statements/ModificationStatement.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/cql3/statements/ModificationStatement.java b/src/java/org/apache/cassandra/cql3/statements/ModificationStatement.java index b960704..6072c24 100644 --- a/src/java/org/apache/cassandra/cql3/statements/ModificationStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/ModificationStatement.java @@ -75,7 +75,14 @@ public abstract class ModificationStatement extends CFStatement implements CQLSt public ResultMessage execute(ClientState state, List<ByteBuffer> variables) throws RequestExecutionException, RequestValidationException { - StorageProxy.mutate(getMutations(state, variables), getConsistencyLevel()); + StorageProxy.mutate(getMutations(state, variables, false), getConsistencyLevel()); + return null; + } + + public ResultMessage executeInternal(ClientState state) throws RequestValidationException, RequestExecutionException + { + for (IMutation mutation : getMutations(state, Collections.<ByteBuffer>emptyList(), true)) + mutation.apply(); return null; } @@ -118,7 +125,7 @@ public abstract class ModificationStatement extends CFStatement implements CQLSt return timeToLive; } - public Map<ByteBuffer, ColumnGroupMap> readRows(List<ByteBuffer> keys, ColumnNameBuilder builder, CompositeType composite) + protected Map<ByteBuffer, ColumnGroupMap> readRows(List<ByteBuffer> keys, ColumnNameBuilder builder, CompositeType composite, boolean local) throws RequestExecutionException, RequestValidationException { List<ReadCommand> commands = new ArrayList<ReadCommand>(keys.size()); @@ -135,7 +142,10 @@ public abstract class ModificationStatement extends CFStatement implements CQLSt try { - List<Row> rows = StorageProxy.read(commands, getConsistencyLevel()); + List<Row> rows = local + ? SelectStatement.readLocally(keyspace(), commands) + : StorageProxy.read(commands, getConsistencyLevel()); + Map<ByteBuffer, ColumnGroupMap> map = new HashMap<ByteBuffer, ColumnGroupMap>(); for (Row row : rows) { @@ -164,11 +174,12 @@ public abstract class ModificationStatement extends CFStatement implements CQLSt * * @param clientState current client status * @param variables value for prepared statement markers + * @param local if true, any requests (for collections) performed by getMutation should be done locally only. * * @return list of the mutations * @throws InvalidRequestException on invalid requests */ - public abstract List<IMutation> getMutations(ClientState clientState, List<ByteBuffer> variables) + protected abstract List<IMutation> getMutations(ClientState clientState, List<ByteBuffer> variables, boolean local) throws RequestExecutionException, RequestValidationException; public abstract ParsedStatement.Prepared prepare(CFDefinition.Name[] boundNames) throws InvalidRequestException; http://git-wip-us.apache.org/repos/asf/cassandra/blob/c74a871e/src/java/org/apache/cassandra/cql3/statements/PermissionAlteringStatement.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/cql3/statements/PermissionAlteringStatement.java b/src/java/org/apache/cassandra/cql3/statements/PermissionAlteringStatement.java new file mode 100644 index 0000000..d3d1c9f --- /dev/null +++ b/src/java/org/apache/cassandra/cql3/statements/PermissionAlteringStatement.java @@ -0,0 +1,48 @@ +/* + * 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.cql3.CQLStatement; +import org.apache.cassandra.service.ClientState; +import org.apache.cassandra.transport.messages.ResultMessage; + +public abstract class PermissionAlteringStatement extends ParsedStatement implements CQLStatement +{ + @Override + public Prepared prepare() + { + return new Prepared(this); + } + + public int getBoundsTerms() + { + return 0; + } + + public void checkAccess(ClientState state) + {} + + public void validate(ClientState state) + {} + + public ResultMessage executeInternal(ClientState state) + { + // executeInternal is for local query only, thus altering permission doesn't make sense and is not supported + throw new UnsupportedOperationException(); + } +} http://git-wip-us.apache.org/repos/asf/cassandra/blob/c74a871e/src/java/org/apache/cassandra/cql3/statements/RevokeStatement.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/cql3/statements/RevokeStatement.java b/src/java/org/apache/cassandra/cql3/statements/RevokeStatement.java index f741642..68c0514 100644 --- a/src/java/org/apache/cassandra/cql3/statements/RevokeStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/RevokeStatement.java @@ -23,13 +23,12 @@ import java.util.List; import org.apache.cassandra.auth.Permission; import org.apache.cassandra.cql3.CFName; -import org.apache.cassandra.cql3.CQLStatement; import org.apache.cassandra.exceptions.InvalidRequestException; import org.apache.cassandra.exceptions.UnauthorizedException; import org.apache.cassandra.service.ClientState; import org.apache.cassandra.transport.messages.ResultMessage; -public class RevokeStatement extends ParsedStatement implements CQLStatement +public class RevokeStatement extends PermissionAlteringStatement { private final Permission permission; private final String from; @@ -42,27 +41,9 @@ public class RevokeStatement extends ParsedStatement implements CQLStatement this.resource = resource; } - public int getBoundsTerms() - { - return 0; - } - - public void checkAccess(ClientState state) throws InvalidRequestException - { - } - - public void validate(ClientState state) throws InvalidRequestException - { - } - public ResultMessage execute(ClientState state, List<ByteBuffer> variables) throws UnauthorizedException, InvalidRequestException { state.revokePermission(permission, from, resource); return null; } - - public Prepared prepare() throws InvalidRequestException - { - return new Prepared(this); - } } http://git-wip-us.apache.org/repos/asf/cassandra/blob/c74a871e/src/java/org/apache/cassandra/cql3/statements/SchemaAlteringStatement.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/cql3/statements/SchemaAlteringStatement.java b/src/java/org/apache/cassandra/cql3/statements/SchemaAlteringStatement.java index 34c957e..36407b4 100644 --- a/src/java/org/apache/cassandra/cql3/statements/SchemaAlteringStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/SchemaAlteringStatement.java @@ -42,8 +42,6 @@ import com.google.common.collect.Maps; */ public abstract class SchemaAlteringStatement extends CFStatement implements CQLStatement { - private static final long timeLimitForSchemaAgreement = 10 * 1000; - private final boolean isColumnFamilyLevel; protected SchemaAlteringStatement() @@ -65,7 +63,8 @@ public abstract class SchemaAlteringStatement extends CFStatement implements CQL super.prepareKeyspace(state); } - public Prepared prepare() throws InvalidRequestException + @Override + public Prepared prepare() { return new Prepared(this); } @@ -90,4 +89,10 @@ public abstract class SchemaAlteringStatement extends CFStatement implements CQL } return null; } + + public ResultMessage executeInternal(ClientState state) + { + // executeInternal is for local query only, thus altering schema is not supported + throw new UnsupportedOperationException(); + } } http://git-wip-us.apache.org/repos/asf/cassandra/blob/c74a871e/src/java/org/apache/cassandra/cql3/statements/SelectStatement.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/cql3/statements/SelectStatement.java b/src/java/org/apache/cassandra/cql3/statements/SelectStatement.java index 9ab0e3c..4af9afd 100644 --- a/src/java/org/apache/cassandra/cql3/statements/SelectStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/SelectStatement.java @@ -20,6 +20,7 @@ package org.apache.cassandra.cql3.statements; import java.io.IOException; import java.nio.ByteBuffer; import java.util.*; +import java.util.concurrent.ExecutionException; import com.google.common.collect.AbstractIterator; import org.slf4j.Logger; @@ -40,6 +41,7 @@ import org.apache.cassandra.db.marshal.*; import org.apache.cassandra.dht.*; import org.apache.cassandra.exceptions.*; import org.apache.cassandra.service.ClientState; +import org.apache.cassandra.service.RangeSliceVerbHandler; import org.apache.cassandra.service.StorageProxy; import org.apache.cassandra.service.StorageService; import org.apache.cassandra.db.ConsistencyLevel; @@ -112,25 +114,59 @@ public class SelectStatement implements CQLStatement public ResultMessage.Rows execute(ClientState state, List<ByteBuffer> variables) throws RequestExecutionException, RequestValidationException { - return new ResultMessage.Rows(executeInternal(state, variables)); - } - - public ResultSet executeInternal(ClientState state, List<ByteBuffer> variables) throws RequestExecutionException, RequestValidationException - { - List<Row> rows; - if (isKeyRange) + try { - rows = multiRangeSlice(variables); + List<Row> rows = isKeyRange + ? StorageProxy.getRangeSlice(getRangeCommand(variables), getConsistencyLevel()) + : StorageProxy.read(getSliceCommands(variables), getConsistencyLevel()); + + return processResults(rows, variables); } - else + catch (IOException e) { - rows = getSlice(variables); + throw new RuntimeException(e); } + } + private ResultMessage.Rows processResults(List<Row> rows, List<ByteBuffer> variables) throws RequestValidationException + { // Even for count, we need to process the result as it'll group some column together in sparse column families ResultSet rset = process(rows, variables); rset = parameters.isCount ? rset.makeCountResult() : rset; - return rset; + return new ResultMessage.Rows(rset); + } + + static List<Row> readLocally(String keyspace, List<ReadCommand> cmds) throws IOException + { + Table table = Table.open(keyspace); + List<Row> rows = new ArrayList(cmds.size()); + for (ReadCommand cmd : cmds) + rows.add(cmd.getRow(table)); + return rows; + } + + public ResultMessage.Rows executeInternal(ClientState state) throws RequestExecutionException, RequestValidationException + { + try + { + List<Row> rows = isKeyRange + ? RangeSliceVerbHandler.executeLocally(getRangeCommand(Collections.<ByteBuffer>emptyList())) + : readLocally(keyspace(), getSliceCommands(Collections.<ByteBuffer>emptyList())); + + return processResults(rows, Collections.<ByteBuffer>emptyList()); + } + catch (IOException e) + { + throw new RuntimeException(e); + } + catch (ExecutionException e) + { + throw new RuntimeException(e); + } + catch (InterruptedException e) + { + throw new RuntimeException(e); + } } public ResultSet process(List<Row> rows) throws InvalidRequestException @@ -149,7 +185,7 @@ public class SelectStatement implements CQLStatement return cfDef.cfm.cfName; } - private List<Row> getSlice(List<ByteBuffer> variables) throws RequestExecutionException, RequestValidationException + private List<ReadCommand> getSliceCommands(List<ByteBuffer> variables) throws RequestValidationException { QueryPath queryPath = new QueryPath(columnFamily()); Collection<ByteBuffer> keys = getKeys(variables); @@ -177,46 +213,27 @@ public class SelectStatement implements CQLStatement commands.add(new SliceByNamesReadCommand(keyspace(), key, queryPath, (NamesQueryFilter)filter)); } } - - try - { - return StorageProxy.read(commands, getConsistencyLevel()); - } - catch (IOException e) - { - throw new RuntimeException(e); - } + return commands; } - private List<Row> multiRangeSlice(List<ByteBuffer> variables) throws RequestExecutionException, RequestValidationException + private RangeSliceCommand getRangeCommand(List<ByteBuffer> variables) throws RequestValidationException { - List<Row> rows; IFilter filter = makeFilter(variables); List<IndexExpression> expressions = getIndexExpressions(variables); - - try - { - // The LIMIT provided by the user is the number of CQL row he wants returned. - // For NamesQueryFilter, this is the number of internal rows returned, since a NamesQueryFilter can only select one CQL row in a given internal row. - // For SliceQueryFilter however, we want to have getRangeSlice to count the number of columns, not the number of keys. Then - // SliceQueryFilter.collectReducedColumns will correctly columns having the same composite prefix using ColumnCounter. - boolean maxIsColumns = filter instanceof SliceQueryFilter; - rows = StorageProxy.getRangeSlice(new RangeSliceCommand(keyspace(), - columnFamily(), - null, - filter, - getKeyBounds(variables), - expressions, - getLimit(), - maxIsColumns, - false), - getConsistencyLevel()); - } - catch (IOException e) - { - throw new RuntimeException(e); - } - return rows; + // The LIMIT provided by the user is the number of CQL row he wants returned. + // For NamesQueryFilter, this is the number of internal rows returned, since a NamesQueryFilter can only select one CQL row in a given internal row. + // For SliceQueryFilter however, we want to have getRangeSlice to count the number of columns, not the number of keys. Then + // SliceQueryFilter.collectReducedColumns will correctly columns having the same composite prefix using ColumnCounter. + boolean maxIsColumns = filter instanceof SliceQueryFilter; + return new RangeSliceCommand(keyspace(), + columnFamily(), + null, + filter, + getKeyBounds(variables), + expressions, + getLimit(), + maxIsColumns, + false); } private AbstractBounds<RowPosition> getKeyBounds(List<ByteBuffer> variables) throws InvalidRequestException http://git-wip-us.apache.org/repos/asf/cassandra/blob/c74a871e/src/java/org/apache/cassandra/cql3/statements/TruncateStatement.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/cql3/statements/TruncateStatement.java b/src/java/org/apache/cassandra/cql3/statements/TruncateStatement.java index f2a1424..e709a06 100644 --- a/src/java/org/apache/cassandra/cql3/statements/TruncateStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/TruncateStatement.java @@ -72,4 +72,9 @@ public class TruncateStatement extends CFStatement implements CQLStatement } return null; } + + public ResultMessage executeInternal(ClientState state) + { + throw new UnsupportedOperationException(); + } } http://git-wip-us.apache.org/repos/asf/cassandra/blob/c74a871e/src/java/org/apache/cassandra/cql3/statements/UpdateStatement.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/cql3/statements/UpdateStatement.java b/src/java/org/apache/cassandra/cql3/statements/UpdateStatement.java index a31d7cb..877d0d8 100644 --- a/src/java/org/apache/cassandra/cql3/statements/UpdateStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/UpdateStatement.java @@ -100,7 +100,7 @@ public class UpdateStatement extends ModificationStatement /** {@inheritDoc} */ - public List<IMutation> getMutations(ClientState clientState, List<ByteBuffer> variables) + public List<IMutation> getMutations(ClientState clientState, List<ByteBuffer> variables, boolean local) throws RequestExecutionException, RequestValidationException { List<ByteBuffer> keys = buildKeyNames(cfDef, processedKeys, variables); @@ -127,7 +127,7 @@ public class UpdateStatement extends ModificationStatement } } - Map<ByteBuffer, ColumnGroupMap> rows = needsReading ? readRows(keys, builder, (CompositeType)cfDef.cfm.comparator) : null; + Map<ByteBuffer, ColumnGroupMap> rows = needsReading ? readRows(keys, builder, (CompositeType)cfDef.cfm.comparator, local) : null; List<IMutation> rowMutations = new LinkedList<IMutation>(); UpdateParameters params = new UpdateParameters(variables, getTimestamp(clientState), getTimeToLive()); http://git-wip-us.apache.org/repos/asf/cassandra/blob/c74a871e/src/java/org/apache/cassandra/cql3/statements/UseStatement.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/cql3/statements/UseStatement.java b/src/java/org/apache/cassandra/cql3/statements/UseStatement.java index f56c6f2..70e95f6 100644 --- a/src/java/org/apache/cassandra/cql3/statements/UseStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/UseStatement.java @@ -53,4 +53,10 @@ public class UseStatement extends ParsedStatement implements CQLStatement state.setKeyspace(keyspace); return new ResultMessage.SetKeyspace(keyspace); } + + public ResultMessage executeInternal(ClientState state) + { + // Internal queries are exclusively on the system keyspace and 'use' is thus useless + throw new UnsupportedOperationException(); + } } http://git-wip-us.apache.org/repos/asf/cassandra/blob/c74a871e/src/java/org/apache/cassandra/service/RangeSliceVerbHandler.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/service/RangeSliceVerbHandler.java b/src/java/org/apache/cassandra/service/RangeSliceVerbHandler.java index e3b1d79..ef057e6 100644 --- a/src/java/org/apache/cassandra/service/RangeSliceVerbHandler.java +++ b/src/java/org/apache/cassandra/service/RangeSliceVerbHandler.java @@ -36,7 +36,7 @@ public class RangeSliceVerbHandler implements IVerbHandler<RangeSliceCommand> { private static final Logger logger = LoggerFactory.getLogger(RangeSliceVerbHandler.class); - static List<Row> executeLocally(RangeSliceCommand command) throws ExecutionException, InterruptedException + public static List<Row> executeLocally(RangeSliceCommand command) throws ExecutionException, InterruptedException { ColumnFamilyStore cfs = Table.open(command.keyspace).getColumnFamilyStore(command.column_family); if (cfs.indexManager.hasIndexFor(command.row_filter))
