From 2e61cd5e07f3983d262ec6bba2aea329e28c5fdc Mon Sep 17 00:00:00 2001 From: Sylvain Lebresne Date: Wed, 7 May 2014 10:53:09 +0200 Subject: [PATCH 1/2] Validate statements inside batch --- .../org/apache/cassandra/cql3/statements/BatchStatement.java | 2 ++ .../cassandra/cql3/statements/ModificationStatement.java | 5 +---- 2 files changed, 3 insertions(+), 4 deletions(-) diff --git a/src/java/org/apache/cassandra/cql3/statements/BatchStatement.java b/src/java/org/apache/cassandra/cql3/statements/BatchStatement.java index c03548bb4e..6a1201b134 100644 --- a/src/java/org/apache/cassandra/cql3/statements/BatchStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/BatchStatement.java @@ -128,6 +128,8 @@ public class BatchStatement implements CQLStatement, MeasurableForPreparedCache { if (timestampSet && statement.isTimestampSet()) throw new InvalidRequestException("Timestamp must be set either on BATCH or individual statements"); + + statement.validate(state); } } diff --git a/src/java/org/apache/cassandra/cql3/statements/ModificationStatement.java b/src/java/org/apache/cassandra/cql3/statements/ModificationStatement.java index 526a26c7fa..f8c40424b7 100644 --- a/src/java/org/apache/cassandra/cql3/statements/ModificationStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/ModificationStatement.java @@ -155,7 +155,7 @@ public abstract class ModificationStatement implements CQLStatement, MeasurableF public void validate(ClientState state) throws InvalidRequestException { if (hasConditions() && attrs.isTimestampSet()) - throw new InvalidRequestException("Custom timestamps are not allowed when conditions are used"); + throw new InvalidRequestException("Cannot provide custom timestamp for conditional update"); if (isCounter()) { @@ -765,9 +765,6 @@ public abstract class ModificationStatement implements CQLStatement, MeasurableF if (stmt.isCounter()) throw new InvalidRequestException("Conditional updates are not supported on counter tables"); - if (attrs.timestamp != null) - throw new InvalidRequestException("Cannot provide custom timestamp for conditional update"); - if (ifNotExists) { // To have both 'IF NOT EXISTS' and some other conditions doesn't make sense. From 05bacaeabc96a6d85fbf908dce8474acffcab730 Mon Sep 17 00:00:00 2001 From: Jason Brown Date: Wed, 7 May 2014 11:58:56 -0700 Subject: [PATCH 2/2] Starting threads in the OutboundTcpConnectionPool constructor causes race conditions patch by sbtourist; reviewed by jasobrown for CASSANDRA-7177 --- CHANGES.txt | 2 +- .../cassandra/net/MessagingService.java | 6 +-- .../net/OutboundTcpConnectionPool.java | 41 ++++++++++++++++--- 3 files changed, 40 insertions(+), 9 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index fc192eff0c..65ee6cf5a9 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,6 +1,6 @@ 2.0.9 * Warn when 'USING TIMESTAMP' is used on a CAS BATCH (CASSANDRA-7067) - + * Starting threads in OutboundTcpConnectionPool constructor causes race conditions (CASSANDRA-7177) 2.0.8 * Correctly delete scheduled range xfers (CASSANDRA-7143) diff --git a/src/java/org/apache/cassandra/net/MessagingService.java b/src/java/org/apache/cassandra/net/MessagingService.java index cccf698346..dbd76d6d99 100644 --- a/src/java/org/apache/cassandra/net/MessagingService.java +++ b/src/java/org/apache/cassandra/net/MessagingService.java @@ -498,11 +498,11 @@ public final class MessagingService implements MessagingServiceMBean cp = new OutboundTcpConnectionPool(to); OutboundTcpConnectionPool existingPool = connectionManagers.putIfAbsent(to, cp); if (existingPool != null) - { - cp.close(); cp = existingPool; - } + else + cp.start(); } + cp.waitForStarted(); return cp; } diff --git a/src/java/org/apache/cassandra/net/OutboundTcpConnectionPool.java b/src/java/org/apache/cassandra/net/OutboundTcpConnectionPool.java index 81168c6873..c45fc530a4 100644 --- a/src/java/org/apache/cassandra/net/OutboundTcpConnectionPool.java +++ b/src/java/org/apache/cassandra/net/OutboundTcpConnectionPool.java @@ -22,6 +22,8 @@ import java.net.InetAddress; import java.net.InetSocketAddress; import java.net.Socket; import java.nio.channels.SocketChannel; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; import org.apache.cassandra.concurrent.Stage; import org.apache.cassandra.config.Config; @@ -36,6 +38,7 @@ public class OutboundTcpConnectionPool { // pointer for the real Address. private final InetAddress id; + private final CountDownLatch started; public final OutboundTcpConnection cmdCon; public final OutboundTcpConnection ackCon; // pointer to the reseted Address. @@ -46,13 +49,10 @@ public class OutboundTcpConnectionPool { id = remoteEp; resetedEndpoint = SystemKeyspace.getPreferredIP(remoteEp); + started = new CountDownLatch(1); cmdCon = new OutboundTcpConnection(this); - cmdCon.start(); ackCon = new OutboundTcpConnection(this); - ackCon.start(); - - metrics = new ConnectionMetrics(id, this); } /** @@ -167,14 +167,45 @@ public class OutboundTcpConnectionPool } return true; } + + public void start() + { + cmdCon.start(); + ackCon.start(); - public void close() + metrics = new ConnectionMetrics(id, this); + + started.countDown(); + } + + public void waitForStarted() + { + if (started.getCount() == 0) + return; + + boolean error = false; + try + { + if (!started.await(1, TimeUnit.MINUTES)) + error = true; + } + catch (InterruptedException e) + { + Thread.currentThread().interrupt(); + error = true; + } + if (error) + throw new IllegalStateException(String.format("Connections to %s are not started!", id.getHostAddress())); + } + + public void close() { // these null guards are simply for tests if (ackCon != null) ackCon.closeSocket(true); if (cmdCon != null) cmdCon.closeSocket(true); + metrics.release(); } }