From c2f24d2c45aae6030310d881dcd96ba60d04a2ad Mon Sep 17 00:00:00 2001 From: Aleksey Yeshchenko Date: Mon, 25 Jan 2021 17:19:56 +0000 Subject: [PATCH] Fix inbound TLS enforcement --- .../org/apache/cassandra/config/Config.java | 1 + .../cassandra/config/DatabaseDescriptor.java | 4 +- .../cassandra/config/EncryptionOptions.java | 48 ++++++ .../net/IncomingStreamingConnection.java | 20 +++ .../cassandra/net/IncomingTcpConnection.java | 22 +++ .../net/OutboundTcpConnectionPool.java | 26 +--- test/conf/cassandra_ssl_test.keystore | Bin 0 -> 2281 bytes test/conf/cassandra_ssl_test.truststore | Bin 0 -> 992 bytes .../InternodeEncryptionEnforcementTest.java | 137 ++++++++++++++++++ 9 files changed, 232 insertions(+), 26 deletions(-) create mode 100644 test/conf/cassandra_ssl_test.keystore create mode 100644 test/conf/cassandra_ssl_test.truststore create mode 100644 test/distributed/org/apache/cassandra/distributed/test/InternodeEncryptionEnforcementTest.java diff --git a/src/java/org/apache/cassandra/config/Config.java b/src/java/org/apache/cassandra/config/Config.java index 277a68a57d..cd74b6140d 100644 --- a/src/java/org/apache/cassandra/config/Config.java +++ b/src/java/org/apache/cassandra/config/Config.java @@ -466,6 +466,7 @@ public class Config private static final List SENSITIVE_KEYS = new ArrayList() {{ add("client_encryption_options"); add("server_encryption_options"); + add("encryption_options"); }}; public static void log(Config config) diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index 52f01b6be7..656828cb16 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -728,13 +728,15 @@ public class DatabaseDescriptor throw new ConfigurationException("index_summary_capacity_in_mb option was set incorrectly to '" + conf.index_summary_capacity_in_mb + "', it should be a non-negative integer.", false); - if(conf.encryption_options != null) + if (conf.encryption_options != null) { logger.warn("Please rename encryption_options as server_encryption_options in the yaml"); //operate under the assumption that server_encryption_options is not set in yaml rather than both conf.server_encryption_options = conf.encryption_options; } + conf.server_encryption_options.validate(); + // load the seeds for node contact points if (conf.seed_provider == null) { diff --git a/src/java/org/apache/cassandra/config/EncryptionOptions.java b/src/java/org/apache/cassandra/config/EncryptionOptions.java index 31f8b4a82d..497768f219 100644 --- a/src/java/org/apache/cassandra/config/EncryptionOptions.java +++ b/src/java/org/apache/cassandra/config/EncryptionOptions.java @@ -17,8 +17,18 @@ */ package org.apache.cassandra.config; +import java.net.InetAddress; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import org.apache.cassandra.locator.IEndpointSnitch; +import org.apache.cassandra.utils.FBUtilities; + public abstract class EncryptionOptions { + private static final Logger logger = LoggerFactory.getLogger(EncryptionOptions.class); + public String keystore = "conf/.keystore"; public String keystore_password = "cassandra"; public String truststore = "conf/.truststore"; @@ -45,6 +55,44 @@ public abstract class EncryptionOptions { all, none, dc, rack } + public InternodeEncryption internode_encryption = InternodeEncryption.none; + + public boolean shouldEncrypt(InetAddress endpoint) + { + IEndpointSnitch snitch = DatabaseDescriptor.getEndpointSnitch(); + InetAddress local = FBUtilities.getBroadcastAddress(); + + switch (internode_encryption) + { + case none: + return false; // if nothing needs to be encrypted then return immediately. + case all: + break; + case dc: + if (snitch.getDatacenter(endpoint).equals(snitch.getDatacenter(local))) + return false; + break; + case rack: + // for rack then check if the DC's are the same. + if (snitch.getRack(endpoint).equals(snitch.getRack(local)) + && snitch.getDatacenter(endpoint).equals(snitch.getDatacenter(local))) + return false; + break; + } + return true; + } + + public void validate() + { + if (require_client_auth && (internode_encryption == InternodeEncryption.rack || internode_encryption == InternodeEncryption.dc)) + { + logger.warn("Setting require_client_auth is incompatible with 'rack' and 'dc' internode_encryption values." + + " It is possible for an internode connection to pretend to be in the same rack/dc by spoofing" + + " its broadcast address in the handshake and bypass authentication. To ensure that mutual TLS" + + " authentication is not bypassed, please set internode_encryption to 'all'. Continuing with" + + " insecure configuration."); + } + } } } diff --git a/src/java/org/apache/cassandra/net/IncomingStreamingConnection.java b/src/java/org/apache/cassandra/net/IncomingStreamingConnection.java index b97b836f1f..77080239d9 100644 --- a/src/java/org/apache/cassandra/net/IncomingStreamingConnection.java +++ b/src/java/org/apache/cassandra/net/IncomingStreamingConnection.java @@ -19,9 +19,12 @@ package org.apache.cassandra.net; import java.io.Closeable; import java.io.IOException; +import java.net.InetAddress; import java.net.Socket; import java.util.Set; +import javax.net.ssl.SSLSocket; + import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -65,6 +68,13 @@ public class IncomingStreamingConnection extends Thread implements Closeable DataInputPlus input = new DataInputStreamPlus(socket.getInputStream()); StreamInitMessage init = StreamInitMessage.serializer.deserialize(input, version); + if (isEncryptionRequired(init.from) && !isEncrypted()) + { + logger.warn("Peer {} attempted to establish an unencrypted streaming connection (broadcast address {})", + socket.getRemoteSocketAddress(), init.from); + throw new IOException("Peer " + init.from + " attempted an unencrypted streaming connection"); + } + //Set SO_TIMEOUT on follower side if (!init.isForOutgoing) socket.setSoTimeout(DatabaseDescriptor.getStreamingSocketTimeout()); @@ -101,4 +111,14 @@ public class IncomingStreamingConnection extends Thread implements Closeable group.remove(this); } } + + private boolean isEncryptionRequired(InetAddress peer) + { + return DatabaseDescriptor.getServerEncryptionOptions().shouldEncrypt(peer); + } + + private boolean isEncrypted() + { + return socket instanceof SSLSocket; + } } diff --git a/src/java/org/apache/cassandra/net/IncomingTcpConnection.java b/src/java/org/apache/cassandra/net/IncomingTcpConnection.java index e79da313ab..ae71e90463 100644 --- a/src/java/org/apache/cassandra/net/IncomingTcpConnection.java +++ b/src/java/org/apache/cassandra/net/IncomingTcpConnection.java @@ -26,6 +26,9 @@ import java.nio.channels.ReadableByteChannel; import java.util.zip.Checksum; import java.util.Set; +import javax.net.ssl.SSLServerSocket; +import javax.net.ssl.SSLSocket; + import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -35,6 +38,7 @@ import net.jpountz.lz4.LZ4Factory; import net.jpountz.xxhash.XXHashFactory; import org.apache.cassandra.config.Config; +import org.apache.cassandra.config.EncryptionOptions; import org.xerial.snappy.SnappyInputStream; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.UnknownColumnFamilyException; @@ -149,6 +153,14 @@ public class IncomingTcpConnection extends Thread implements Closeable DataInputPlus in = new DataInputStreamPlus(socket.getInputStream()); int maxVersion = in.readInt(); from = CompactEndpointSerializationHelper.deserialize(in); + + if (isEncryptionRequired(from) && !isEncrypted()) + { + logger.warn("Peer {} attempted to establish an unencrypted connection (broadcast address {})", + socket.getRemoteSocketAddress(), from); + throw new IOException("Peer " + from + " attempted an unencrypted connection"); + } + // record the (true) version of the endpoint MessagingService.instance().setVersion(from, maxVersion); logger.trace("Set version for {} to {} (will use {})", from, maxVersion, MessagingService.instance().getVersion(from)); @@ -217,4 +229,14 @@ public class IncomingTcpConnection extends Thread implements Closeable } return message.from; } + + private boolean isEncryptionRequired(InetAddress peer) + { + return DatabaseDescriptor.getServerEncryptionOptions().shouldEncrypt(peer); + } + + private boolean isEncrypted() + { + return socket instanceof SSLSocket; + } } diff --git a/src/java/org/apache/cassandra/net/OutboundTcpConnectionPool.java b/src/java/org/apache/cassandra/net/OutboundTcpConnectionPool.java index 2b92036970..c9dfe9ca90 100644 --- a/src/java/org/apache/cassandra/net/OutboundTcpConnectionPool.java +++ b/src/java/org/apache/cassandra/net/OutboundTcpConnectionPool.java @@ -29,7 +29,6 @@ import org.apache.cassandra.concurrent.Stage; import org.apache.cassandra.config.Config; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.SystemKeyspace; -import org.apache.cassandra.locator.IEndpointSnitch; import org.apache.cassandra.metrics.ConnectionMetrics; import org.apache.cassandra.security.SSLFactory; import org.apache.cassandra.utils.FBUtilities; @@ -127,7 +126,7 @@ public class OutboundTcpConnectionPool public static Socket newSocket(InetAddress endpoint) throws IOException { // zero means 'bind on any available port.' - if (isEncryptedChannel(endpoint)) + if (DatabaseDescriptor.getServerEncryptionOptions().shouldEncrypt(endpoint)) { if (DatabaseDescriptor.getOutboundBindAny()) return SSLFactory.getSocket(DatabaseDescriptor.getServerEncryptionOptions(), endpoint, DatabaseDescriptor.getSSLStoragePort()); @@ -151,29 +150,6 @@ public class OutboundTcpConnectionPool return resetEndpoint; } - public static boolean isEncryptedChannel(InetAddress address) - { - IEndpointSnitch snitch = DatabaseDescriptor.getEndpointSnitch(); - switch (DatabaseDescriptor.getServerEncryptionOptions().internode_encryption) - { - case none: - return false; // if nothing needs to be encrypted then return immediately. - case all: - break; - case dc: - if (snitch.getDatacenter(address).equals(snitch.getDatacenter(FBUtilities.getBroadcastAddress()))) - return false; - break; - case rack: - // for rack then check if the DC's are the same. - if (snitch.getRack(address).equals(snitch.getRack(FBUtilities.getBroadcastAddress())) - && snitch.getDatacenter(address).equals(snitch.getDatacenter(FBUtilities.getBroadcastAddress()))) - return false; - break; - } - return true; - } - public void start() { smallMessages.start(); diff --git a/test/conf/cassandra_ssl_test.keystore b/test/conf/cassandra_ssl_test.keystore new file mode 100644 index 0000000000000000000000000000000000000000..8b2b218efab60b26583bc620b281941f10dc0b5c GIT binary patch literal 2281 zcmc(g`8N~{7sqE~3S-|HOEE+-V_y4`X|l$U6&iucz~N-sha>4|spLKitnf=bn4M_nz+3 z2O{m(merQuN|(%qmt$%7epq#|Ah~ zus$ocl=|Mt6S1UF*`_`IbJOTLDqzap$NLEM#Toj=k7pGS_IUjpZ@pU`VQoy#YlnZaH8k`Gd`N4gA$1C?!u4ptFZk8KW=1NHXf{xVCMF%k~pgvstJ z&Mp1EE|AK?&diz-*8W+LBf+wubP-Fg@~q`$&1WRdeq}N%MJ;cL1by$&g8e|-oWe()-&aty!>5fPA&lVs0q4D&r6 zEx{aZ*L!7O;5ksYfBH*;>HBccQweiS$Ne^@54uG{PgjfEMXGgSV<2L*)yhcpNe~}I z^ynvsTpC-AG*;ehAA9$VL!-^X3C8{Cx@BBZfvdkuj&GJ{j6ZkP1uaq6uNDkIMkf9_ zZ3|^L&^zkE;W~}@sTWuQP|KHXv?tikFLtV#YcZm#hS;LFyt>G*n>I3YKkz(bnYbgK z!JH*xdhYZCrykyKhGniNr6tKhJ1yacv-*`nl58GJc2^?ZYN_>(wEA>?ZV}7c;4P{2 z4B~qT>q zI>cGh#t@gK*T+QJwJQDlX*6R#J>Z;Q#n1kUq73btxZDAdLbN^v)G6y(phs|yTuj;; z{KQWWV#zgQbpJsbO#k?~djxVq&`Vh?GbW`tZ9_seYY1IBuef9_BJV znHN$Kaa(qHpM=p#y=G(;)2&oKke;LL%ez$9jhUm(9Q5%MFKI0(e(1Hm2nboyQJhod zZ(XE450$EeXixjOw`){&_=-#FHh@#$)f|P#pz@(o^OVZLN(Y}b(?M7=L3Di~13go9 z=jYuz>QH>WOEch)MH47G?#n_eS{_OU@#&fW*sdwU=QI>1Xf7H(3Xt+y8i0I5MD`8+ zp}%vF0$^>*U9M0r*L5-^^;xYv!}Qt&v6lO}M^D5i>*BpSK=JM+($d=x(GADB=GWTz zJ7Y=JbE!7{W!Mj(wu{t2f6(JO@w4~^rwQGg4co3vNX^XSD;zVsE#Y{eMR9>KR`2uu z^7EKAx-`*b^h^5DdMqTk91F7x7hn_$%5NCU4hBAnV~Iu`<118~ZR+-ZrBN(+N#qnR z_Ci`Al^rPw)0-tvmNfL*h+Vh7VdO8yKN+;4R<{LEYIi2{n#GkmoO4b*-<`5^pVsut zBsrYaG?BUh05}s#1>ZqZK@s_25D)}{*a-YoK*BhWs?f}7-(f%?HwOS2f#N@OOt2$_ zQwZuvCXoZe$iK;vzY*x)2-p9JFjDZ(KHNOqzpdGq?2VE@iv3B!Lj0IOf|n0b4*T1| zkuqX1lok?=($PR^=pZ#N{z|n_TK`S|pQBNM$NnnR;Tr){fqVcY703mm0)c>3i%j>m z;n^|4)rjpPa)K22NT5krSRFwYlR}YIr7O6OnK^8_(HkK}0l;`v$33hT>>BzS$xt+) z%^a&ZWT;sb&6#DRX!_i=6Bp^Rv@74YetOGb{Yl={*8#W&gl3q;gDL!GNghrgHt+kQ zxLS{QLc1!m9975E-4>Fe#g`aKJ0LI+04$P6$|B*1 z+lBK&L?Oag$8Tyss=k*Nk`4`A9D`j()BLX zdeMHR1)0cG-T6+~ zS!EyHx%V!s$vm~p;zia$6ejv>N@5D_sD0-3|j!wqv$W?s6Xq=7g{lv!B7u^=%yBUQl}=5PZ!ab80+ z17kxABSRw#1EVM~*UZoi${k3jH8Cn72NWYK19KB2KZ8LNBNtN>BO^nf=hXO*mmXgi z`cV06abAZsGe?1IdHJ?PdB@&jdF{C>vDe*0zQoPl#j-e`vDNVC6lXKuD&wkbTe18@ z9?q&4ZA}-~u}$+=bzARt%&#Wl)gQ&PpX|Tcd|ExfaDTqvPBqhV$r<;5n50fIeu*&D zkX7Gvar@Em{Wl-<)^6FsyX^Aa-hx`LnXfbavZn=I+9UjHj`>z@^D}`FSJfX$@0yXf zEL+neE~3qCk!k4cm8;}?8g8|ptDao$b|7hsj>sQd*=1#IU(;2;pL`bgyoY^B%H{6g zOw5c7jEfZwKSlF`Qq-@>hgY_r0xwmHq9@}u0{mJeRC*L{0H(8V={T zbXf@Al6gAi(BsQhS37o0=(6AZV`l2wnVO6B^&3U&aw;yWH`QJ|Gq?NG>&dG8e?Gj{ iYRGywzgSfK&JK;;RtxssiVdn*{;X0_$}!}W+5`ZGs(6b4 literal 0 HcmV?d00001 diff --git a/test/distributed/org/apache/cassandra/distributed/test/InternodeEncryptionEnforcementTest.java b/test/distributed/org/apache/cassandra/distributed/test/InternodeEncryptionEnforcementTest.java new file mode 100644 index 0000000000..5c0e3b3f93 --- /dev/null +++ b/test/distributed/org/apache/cassandra/distributed/test/InternodeEncryptionEnforcementTest.java @@ -0,0 +1,137 @@ +/* + * 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.distributed.test; + +import java.util.HashMap; +import java.util.List; + +import com.google.common.collect.ImmutableMap; +import org.junit.Test; + +import org.apache.cassandra.distributed.Cluster; +import org.apache.cassandra.distributed.api.Feature; +import org.apache.cassandra.distributed.api.IIsolatedExecutor.SerializableRunnable; +import org.apache.cassandra.distributed.shared.NetworkTopology; +import org.apache.cassandra.net.MessagingService; + +import static com.google.common.collect.Iterables.getOnlyElement; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; + +public final class InternodeEncryptionEnforcementTest extends TestBaseImpl +{ + @Test + public void testConnectionsAreRejectedWithInvalidConfig() throws Throwable + { + Cluster.Builder builder = builder() + .withNodes(2) + .withConfig(c -> + { + c.with(Feature.NETWORK); + c.with(Feature.NATIVE_PROTOCOL); + + if (c.num() == 1) + { + HashMap encryption = new HashMap<>(); + encryption.put("keystore", "test/conf/cassandra_ssl_test.keystore"); + encryption.put("keystore_password", "cassandra"); + encryption.put("truststore", "test/conf/cassandra_ssl_test.truststore"); + encryption.put("truststore_password", "cassandra"); + encryption.put("internode_encryption", "dc"); + c.set("server_encryption_options", encryption); + } + }) + .withNodeIdTopology(ImmutableMap.of(1, NetworkTopology.dcAndRack("dc1", "r1a"), + 2, NetworkTopology.dcAndRack("dc2", "r2a"))); + + try (Cluster cluster = builder.start()) + { + /* + * instance (1) won't connect to (2), since (2) won't have a TLS listener; + * instance (2) won't connect to (1), since inbound check will reject + * the unencrypted connection attempt; + * + * without the patch, instance (2) *CAN* connect to (1), without encryption, + * despite being in a different dc. + */ + + cluster.get(1).runOnInstance(() -> + { + List threads = MessagingService.instance().getSocketThreads(); + assertEquals(2, threads.size()); + + for (MessagingService.SocketThread thread : threads) + { + assertEquals(0, thread.connections.size()); + } + }); + + cluster.get(2).runOnInstance(() -> + { + List threads = MessagingService.instance().getSocketThreads(); + assertEquals(1, threads.size()); + assertTrue(getOnlyElement(threads).connections.isEmpty()); + }); + } + } + + @Test + public void testConnectionsAreAcceptedWithValidConfig() throws Throwable + { + Cluster.Builder builder = builder() + .withNodes(2) + .withConfig(c -> + { + c.with(Feature.NETWORK); + c.with(Feature.NATIVE_PROTOCOL); + + HashMap encryption = new HashMap<>(); + encryption.put("keystore", "test/conf/cassandra_ssl_test.keystore"); + encryption.put("keystore_password", "cassandra"); + encryption.put("truststore", "test/conf/cassandra_ssl_test.truststore"); + encryption.put("truststore_password", "cassandra"); + encryption.put("internode_encryption", "dc"); + c.set("server_encryption_options", encryption); + }) + .withNodeIdTopology(ImmutableMap.of(1, NetworkTopology.dcAndRack("dc1", "r1a"), + 2, NetworkTopology.dcAndRack("dc2", "r2a"))); + + try (Cluster cluster = builder.start()) + { + /* + * instance (1) should connect to instance (2) without any issues; + * instance (2) should connect to instance (1) without any issues. + */ + + SerializableRunnable runnable = () -> + { + List threads = MessagingService.instance().getSocketThreads(); + assertEquals(2, threads.size()); + + MessagingService.SocketThread sslThread = threads.get(0); + assertEquals(1, sslThread.connections.size()); + + MessagingService.SocketThread plainThread = threads.get(1); + assertEquals(0, plainThread.connections.size()); + }; + + cluster.get(1).runOnInstance(runnable); + cluster.get(2).runOnInstance(runnable); + } + } +}