diff --git a/CHANGES.txt b/CHANGES.txt
index 5c9b5fadfc..ea5d6cdacc 100644
--- a/CHANGES.txt
+++ b/CHANGES.txt
@@ -11,6 +11,11 @@
* move gossip heartbeat back to its own thread (CASSANDRA-2554)
* validate cql TRUNCATE columnfamily before truncating (CASSANDRA-2570)
* fix batch_mutate for mixed standard-counter mutations (CASSANDRA-2457)
+ * disallow making schema changes to system keyspace (CASSANDRA-2563)
+ * fix sending mutation messages multiple times (CASSANDRA-2557)
+ * fix incorrect use of NBHM.size in ReadCallback that could cause
+ reads to time out even when responses were received (CASSAMDRA-2552)
+ * trigger read repair correctly for LOCAL_QUORUM reads (CASSANDRA-2556)
0.8.0-beta1
diff --git a/build.xml b/build.xml
index 1ee7bba351..f2d229f420 100644
--- a/build.xml
+++ b/build.xml
@@ -618,7 +618,7 @@
-
+
diff --git a/conf/cassandra-env.sh b/conf/cassandra-env.sh
index 16c29efdd9..f6e41d8459 100644
--- a/conf/cassandra-env.sh
+++ b/conf/cassandra-env.sh
@@ -95,7 +95,7 @@ JVM_OPTS="$JVM_OPTS -ea"
check_openjdk=$(java -version 2>&1 | awk '{if (NR == 2) {print $1}}')
if [ "$check_openjdk" != "OpenJDK" ]
then
- JVM_OPTS="$JVM_OPTS -javaagent:$CASSANDRA_HOME/lib/jamm-0.2.1.jar"
+ JVM_OPTS="$JVM_OPTS -javaagent:$CASSANDRA_HOME/lib/jamm-0.2.2.jar"
fi
# enable thread priorities, primarily so we can give periodic tasks
diff --git a/lib/jamm-0.2.1.jar b/lib/jamm-0.2.1.jar
deleted file mode 100644
index 16dd22a624..0000000000
Binary files a/lib/jamm-0.2.1.jar and /dev/null differ
diff --git a/lib/jamm-0.2.2.jar b/lib/jamm-0.2.2.jar
new file mode 100644
index 0000000000..ce4227a950
Binary files /dev/null and b/lib/jamm-0.2.2.jar differ
diff --git a/lib/licenses/jamm-0.2.1.txt b/lib/licenses/jamm-0.2.2.txt
similarity index 100%
rename from lib/licenses/jamm-0.2.1.txt
rename to lib/licenses/jamm-0.2.2.txt
diff --git a/src/java/org/apache/cassandra/service/AbstractRowResolver.java b/src/java/org/apache/cassandra/service/AbstractRowResolver.java
index a6a9a1e3e3..18d44e4f0a 100644
--- a/src/java/org/apache/cassandra/service/AbstractRowResolver.java
+++ b/src/java/org/apache/cassandra/service/AbstractRowResolver.java
@@ -83,9 +83,4 @@ public abstract class AbstractRowResolver implements IResponseResolver
{
return replies.keySet();
}
-
- public int getMessageCount()
- {
- return replies.size();
- }
}
diff --git a/src/java/org/apache/cassandra/service/AsyncRepairCallback.java b/src/java/org/apache/cassandra/service/AsyncRepairCallback.java
index 6c925877e4..72a4b0929b 100644
--- a/src/java/org/apache/cassandra/service/AsyncRepairCallback.java
+++ b/src/java/org/apache/cassandra/service/AsyncRepairCallback.java
@@ -22,6 +22,7 @@ package org.apache.cassandra.service;
import java.io.IOException;
+import java.util.concurrent.atomic.AtomicInteger;
import org.apache.cassandra.concurrent.Stage;
import org.apache.cassandra.concurrent.StageManager;
@@ -32,18 +33,19 @@ import org.apache.cassandra.utils.WrappedRunnable;
public class AsyncRepairCallback implements IAsyncCallback
{
private final RowRepairResolver repairResolver;
- private final int count;
+ private final int blockfor;
+ protected final AtomicInteger received = new AtomicInteger(0);
- public AsyncRepairCallback(RowRepairResolver repairResolver, int count)
+ public AsyncRepairCallback(RowRepairResolver repairResolver, int blockfor)
{
this.repairResolver = repairResolver;
- this.count = count;
+ this.blockfor = blockfor;
}
public void response(Message message)
{
repairResolver.preprocess(message);
- if (repairResolver.getMessageCount() == count)
+ if (received.incrementAndGet() == blockfor)
{
StageManager.getStage(Stage.READ_REPAIR).execute(new WrappedRunnable()
{
diff --git a/src/java/org/apache/cassandra/service/ClientState.java b/src/java/org/apache/cassandra/service/ClientState.java
index 376b909671..c274f763cc 100644
--- a/src/java/org/apache/cassandra/service/ClientState.java
+++ b/src/java/org/apache/cassandra/service/ClientState.java
@@ -129,7 +129,11 @@ public class ClientState
{
validateLogin();
validateKeyspace();
-
+
+ // hardcode disallowing messing with system keyspace
+ if (keyspace.equalsIgnoreCase("system"))
+ throw new InvalidRequestException("system keyspace is not user-modifiable");
+
resourceClear();
resource.add(keyspace);
Set perms = DatabaseDescriptor.getAuthority().authorize(user, resource);
diff --git a/src/java/org/apache/cassandra/service/DatacenterReadCallback.java b/src/java/org/apache/cassandra/service/DatacenterReadCallback.java
index 9cdfa19c64..7d03aa4b67 100644
--- a/src/java/org/apache/cassandra/service/DatacenterReadCallback.java
+++ b/src/java/org/apache/cassandra/service/DatacenterReadCallback.java
@@ -23,7 +23,6 @@ package org.apache.cassandra.service;
import java.net.InetAddress;
import java.util.List;
-import java.util.concurrent.atomic.AtomicInteger;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.ReadResponse;
@@ -42,43 +41,26 @@ public class DatacenterReadCallback extends ReadCallback
{
private static final IEndpointSnitch snitch = DatabaseDescriptor.getEndpointSnitch();
private static final String localdc = snitch.getDatacenter(FBUtilities.getLocalAddress());
- private AtomicInteger localResponses;
-
+
public DatacenterReadCallback(IResponseResolver resolver, ConsistencyLevel consistencyLevel, IReadCommand command, List endpoints)
{
super(resolver, consistencyLevel, command, endpoints);
- localResponses = new AtomicInteger(blockfor);
}
@Override
- public void response(Message message)
+ protected boolean waitingFor(Message message)
{
- resolver.preprocess(message);
-
- int n = localdc.equals(snitch.getDatacenter(message.getFrom()))
- ? localResponses.decrementAndGet()
- : localResponses.get();
-
- if (n == 0 && resolver.isDataPresent())
- {
- condition.signal();
- }
+ return localdc.equals(snitch.getDatacenter(message.getFrom()));
}
-
+
@Override
- public void response(ReadResponse result)
+ protected boolean waitingFor(ReadResponse response)
{
- ((RowDigestResolver) resolver).injectPreProcessed(result);
-
- int n = localResponses.decrementAndGet();
- if (n == 0 && resolver.isDataPresent())
- {
- condition.signal();
- }
-
- maybeResolveForRepair();
+ // cheat and leverage our knowledge that a local read is the only way the ReadResponse
+ // version of this method gets called
+ return true;
}
-
+
@Override
public int determineBlockFor(ConsistencyLevel consistency_level, String table)
{
diff --git a/src/java/org/apache/cassandra/service/IResponseResolver.java b/src/java/org/apache/cassandra/service/IResponseResolver.java
index 59b9bfd44e..f4f972a304 100644
--- a/src/java/org/apache/cassandra/service/IResponseResolver.java
+++ b/src/java/org/apache/cassandra/service/IResponseResolver.java
@@ -43,5 +43,4 @@ public interface IResponseResolver {
public void preprocess(Message message);
public Iterable getMessages();
- public int getMessageCount();
}
diff --git a/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java b/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java
index 521035c9d5..6077e38bc7 100644
--- a/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java
+++ b/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java
@@ -145,9 +145,4 @@ public class RangeSliceResponseResolver implements IResponseResolver implements IAsyncCallback
protected final int blockfor;
final List endpoints;
private final IReadCommand command;
+ protected final AtomicInteger received = new AtomicInteger(0);
/**
* Constructor when response count has to be calculated and blocked for.
@@ -115,7 +117,7 @@ public class ReadCallback implements IAsyncCallback
StringBuilder sb = new StringBuilder("");
for (Message message : resolver.getMessages())
sb.append(message.getFrom()).append(", ");
- throw new TimeoutException("Operation timed out - received only " + resolver.getMessageCount() + " responses from " + sb.toString() + " .");
+ throw new TimeoutException("Operation timed out - received only " + received.get() + " responses from " + sb.toString() + " .");
}
return blockfor == 1 ? resolver.getData() : resolver.resolve();
@@ -124,23 +126,40 @@ public class ReadCallback implements IAsyncCallback
public void response(Message message)
{
resolver.preprocess(message);
- assert resolver.getMessageCount() <= endpoints.size();
- if (resolver.getMessageCount() < blockfor)
- return;
- if (resolver.isDataPresent())
+ int n = waitingFor(message)
+ ? received.incrementAndGet()
+ : received.get();
+ if (n >= blockfor && resolver.isDataPresent())
{
condition.signal();
maybeResolveForRepair();
}
}
+ /**
+ * @return true if the message counts towards the blockfor threshold
+ * TODO turn the Message into a response so we don't need two versions of this method
+ */
+ protected boolean waitingFor(Message message)
+ {
+ return true;
+ }
+
+ /**
+ * @return true if the response counts towards the blockfor threshold
+ */
+ protected boolean waitingFor(ReadResponse response)
+ {
+ return true;
+ }
+
public void response(ReadResponse result)
{
((RowDigestResolver) resolver).injectPreProcessed(result);
- assert resolver.getMessageCount() <= endpoints.size();
- if (resolver.getMessageCount() < blockfor)
- return;
- if (resolver.isDataPresent())
+ int n = waitingFor(result)
+ ? received.incrementAndGet()
+ : received.get();
+ if (n >= blockfor && resolver.isDataPresent())
{
condition.signal();
maybeResolveForRepair();
@@ -153,7 +172,7 @@ public class ReadCallback implements IAsyncCallback
*/
protected void maybeResolveForRepair()
{
- if (blockfor < endpoints.size() && resolver.getMessageCount() == endpoints.size())
+ if (blockfor < endpoints.size() && received.get() == endpoints.size())
{
assert resolver.isDataPresent();
StageManager.getStage(Stage.READ_REPAIR).execute(new AsyncRepairRunner());
diff --git a/src/java/org/apache/cassandra/service/RepairCallback.java b/src/java/org/apache/cassandra/service/RepairCallback.java
index 2b946223a9..d79ea1dcd6 100644
--- a/src/java/org/apache/cassandra/service/RepairCallback.java
+++ b/src/java/org/apache/cassandra/service/RepairCallback.java
@@ -26,6 +26,7 @@ import java.net.InetAddress;
import java.util.List;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
+import java.util.concurrent.atomic.AtomicInteger;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.net.IAsyncCallback;
@@ -38,6 +39,7 @@ public class RepairCallback implements IAsyncCallback
private final List endpoints;
private final SimpleCondition condition = new SimpleCondition();
private final long startTime;
+ protected final AtomicInteger received = new AtomicInteger(0);
/**
* The main difference between this and ReadCallback is, ReadCallback has a ConsistencyLevel
@@ -66,13 +68,13 @@ public class RepairCallback implements IAsyncCallback
throw new AssertionError(ex);
}
- return resolver.getMessageCount() > 1 ? resolver.resolve() : null;
+ return received.get() > 1 ? resolver.resolve() : null;
}
public void response(Message message)
{
resolver.preprocess(message);
- if (resolver.getMessageCount() == endpoints.size())
+ if (received.incrementAndGet() == endpoints.size())
condition.signal();
}
diff --git a/src/resources/org/apache/cassandra/cli/CliHelp.yaml b/src/resources/org/apache/cassandra/cli/CliHelp.yaml
index 2b927f0162..64dcfaa205 100644
--- a/src/resources/org/apache/cassandra/cli/CliHelp.yaml
+++ b/src/resources/org/apache/cassandra/cli/CliHelp.yaml
@@ -426,7 +426,7 @@ commands:
store the whole values of its rows, so it is extremely space-intensive.
It's best to only use the row cache if you have hot rows or static rows.
- - keys_cached_save_period: Duration in seconds after which Cassandra should
+ - keys_cache_save_period: Duration in seconds after which Cassandra should
safe the keys cache. Caches are saved to saved_caches_directory as
specified in conf/Cassandra.yaml. Default is 14400 or 4 hours.
@@ -676,7 +676,7 @@ commands:
store the whole values of its rows, so it is extremely space-intensive.
It's best to only use the row cache if you have hot rows or static rows.
- - keys_cached_save_period: Duration in seconds after which Cassandra should
+ - keys_cache_save_period: Duration in seconds after which Cassandra should
safe the keys cache. Caches are saved to saved_caches_directory as
specified in conf/Cassandra.yaml. Default is 14400 or 4 hours.