diff --git a/CHANGES.txt b/CHANGES.txt
index 2faca24e20..4d3a9a9051 100644
--- a/CHANGES.txt
+++ b/CHANGES.txt
@@ -1,4 +1,5 @@
4.0
+ * Update Netty dependencies to latest, clean up SocketFactory (CASSANDRA-15195)
* Native Transport - Apply noSpamLogger to ConnectionLimitHandler (CASSANDRA-15167)
* Reduce heap pressure during compactions (CASSANDRA-14654)
* Support building Cassandra with JDK 11 (CASSANDRA-15108)
diff --git a/build.xml b/build.xml
index bdf5ae2ab0..acfc6133c7 100644
--- a/build.xml
+++ b/build.xml
@@ -548,7 +548,8 @@
-
+
+
diff --git a/lib/licenses/netty-4.1.28.txt b/lib/licenses/netty-4.1.37.txt
similarity index 100%
rename from lib/licenses/netty-4.1.28.txt
rename to lib/licenses/netty-4.1.37.txt
diff --git a/lib/licenses/netty-tcnative-2.0.25.txt b/lib/licenses/netty-tcnative-2.0.25.txt
new file mode 100644
index 0000000000..261eeb9e9f
--- /dev/null
+++ b/lib/licenses/netty-tcnative-2.0.25.txt
@@ -0,0 +1,201 @@
+ Apache License
+ Version 2.0, January 2004
+ http://www.apache.org/licenses/
+
+ TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION
+
+ 1. Definitions.
+
+ "License" shall mean the terms and conditions for use, reproduction,
+ and distribution as defined by Sections 1 through 9 of this document.
+
+ "Licensor" shall mean the copyright owner or entity authorized by
+ the copyright owner that is granting the License.
+
+ "Legal Entity" shall mean the union of the acting entity and all
+ other entities that control, are controlled by, or are under common
+ control with that entity. For the purposes of this definition,
+ "control" means (i) the power, direct or indirect, to cause the
+ direction or management of such entity, whether by contract or
+ otherwise, or (ii) ownership of fifty percent (50%) or more of the
+ outstanding shares, or (iii) beneficial ownership of such entity.
+
+ "You" (or "Your") shall mean an individual or Legal Entity
+ exercising permissions granted by this License.
+
+ "Source" form shall mean the preferred form for making modifications,
+ including but not limited to software source code, documentation
+ source, and configuration files.
+
+ "Object" form shall mean any form resulting from mechanical
+ transformation or translation of a Source form, including but
+ not limited to compiled object code, generated documentation,
+ and conversions to other media types.
+
+ "Work" shall mean the work of authorship, whether in Source or
+ Object form, made available under the License, as indicated by a
+ copyright notice that is included in or attached to the work
+ (an example is provided in the Appendix below).
+
+ "Derivative Works" shall mean any work, whether in Source or Object
+ form, that is based on (or derived from) the Work and for which the
+ editorial revisions, annotations, elaborations, or other modifications
+ represent, as a whole, an original work of authorship. For the purposes
+ of this License, Derivative Works shall not include works that remain
+ separable from, or merely link (or bind by name) to the interfaces of,
+ the Work and Derivative Works thereof.
+
+ "Contribution" shall mean any work of authorship, including
+ the original version of the Work and any modifications or additions
+ to that Work or Derivative Works thereof, that is intentionally
+ submitted to Licensor for inclusion in the Work by the copyright owner
+ or by an individual or Legal Entity authorized to submit on behalf of
+ the copyright owner. For the purposes of this definition, "submitted"
+ means any form of electronic, verbal, or written communication sent
+ to the Licensor or its representatives, including but not limited to
+ communication on electronic mailing lists, source code control systems,
+ and issue tracking systems that are managed by, or on behalf of, the
+ Licensor for the purpose of discussing and improving the Work, but
+ excluding communication that is conspicuously marked or otherwise
+ designated in writing by the copyright owner as "Not a Contribution."
+
+ "Contributor" shall mean Licensor and any individual or Legal Entity
+ on behalf of whom a Contribution has been received by Licensor and
+ subsequently incorporated within the Work.
+
+ 2. Grant of Copyright License. Subject to the terms and conditions of
+ this License, each Contributor hereby grants to You a perpetual,
+ worldwide, non-exclusive, no-charge, royalty-free, irrevocable
+ copyright license to reproduce, prepare Derivative Works of,
+ publicly display, publicly perform, sublicense, and distribute the
+ Work and such Derivative Works in Source or Object form.
+
+ 3. Grant of Patent License. Subject to the terms and conditions of
+ this License, each Contributor hereby grants to You a perpetual,
+ worldwide, non-exclusive, no-charge, royalty-free, irrevocable
+ (except as stated in this section) patent license to make, have made,
+ use, offer to sell, sell, import, and otherwise transfer the Work,
+ where such license applies only to those patent claims licensable
+ by such Contributor that are necessarily infringed by their
+ Contribution(s) alone or by combination of their Contribution(s)
+ with the Work to which such Contribution(s) was submitted. If You
+ institute patent litigation against any entity (including a
+ cross-claim or counterclaim in a lawsuit) alleging that the Work
+ or a Contribution incorporated within the Work constitutes direct
+ or contributory patent infringement, then any patent licenses
+ granted to You under this License for that Work shall terminate
+ as of the date such litigation is filed.
+
+ 4. Redistribution. You may reproduce and distribute copies of the
+ Work or Derivative Works thereof in any medium, with or without
+ modifications, and in Source or Object form, provided that You
+ meet the following conditions:
+
+ (a) You must give any other recipients of the Work or
+ Derivative Works a copy of this License; and
+
+ (b) You must cause any modified files to carry prominent notices
+ stating that You changed the files; and
+
+ (c) You must retain, in the Source form of any Derivative Works
+ that You distribute, all copyright, patent, trademark, and
+ attribution notices from the Source form of the Work,
+ excluding those notices that do not pertain to any part of
+ the Derivative Works; and
+
+ (d) If the Work includes a "NOTICE" text file as part of its
+ distribution, then any Derivative Works that You distribute must
+ include a readable copy of the attribution notices contained
+ within such NOTICE file, excluding those notices that do not
+ pertain to any part of the Derivative Works, in at least one
+ of the following places: within a NOTICE text file distributed
+ as part of the Derivative Works; within the Source form or
+ documentation, if provided along with the Derivative Works; or,
+ within a display generated by the Derivative Works, if and
+ wherever such third-party notices normally appear. The contents
+ of the NOTICE file are for informational purposes only and
+ do not modify the License. You may add Your own attribution
+ notices within Derivative Works that You distribute, alongside
+ or as an addendum to the NOTICE text from the Work, provided
+ that such additional attribution notices cannot be construed
+ as modifying the License.
+
+ You may add Your own copyright statement to Your modifications and
+ may provide additional or different license terms and conditions
+ for use, reproduction, or distribution of Your modifications, or
+ for any such Derivative Works as a whole, provided Your use,
+ reproduction, and distribution of the Work otherwise complies with
+ the conditions stated in this License.
+
+ 5. Submission of Contributions. Unless You explicitly state otherwise,
+ any Contribution intentionally submitted for inclusion in the Work
+ by You to the Licensor shall be under the terms and conditions of
+ this License, without any additional terms or conditions.
+ Notwithstanding the above, nothing herein shall supersede or modify
+ the terms of any separate license agreement you may have executed
+ with Licensor regarding such Contributions.
+
+ 6. Trademarks. This License does not grant permission to use the trade
+ names, trademarks, service marks, or product names of the Licensor,
+ except as required for reasonable and customary use in describing the
+ origin of the Work and reproducing the content of the NOTICE file.
+
+ 7. Disclaimer of Warranty. Unless required by applicable law or
+ agreed to in writing, Licensor provides the Work (and each
+ Contributor provides its Contributions) on an "AS IS" BASIS,
+ WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or
+ implied, including, without limitation, any warranties or conditions
+ of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A
+ PARTICULAR PURPOSE. You are solely responsible for determining the
+ appropriateness of using or redistributing the Work and assume any
+ risks associated with Your exercise of permissions under this License.
+
+ 8. Limitation of Liability. In no event and under no legal theory,
+ whether in tort (including negligence), contract, or otherwise,
+ unless required by applicable law (such as deliberate and grossly
+ negligent acts) or agreed to in writing, shall any Contributor be
+ liable to You for damages, including any direct, indirect, special,
+ incidental, or consequential damages of any character arising as a
+ result of this License or out of the use or inability to use the
+ Work (including but not limited to damages for loss of goodwill,
+ work stoppage, computer failure or malfunction, or any and all
+ other commercial damages or losses), even if such Contributor
+ has been advised of the possibility of such damages.
+
+ 9. Accepting Warranty or Additional Liability. While redistributing
+ the Work or Derivative Works thereof, You may choose to offer,
+ and charge a fee for, acceptance of support, warranty, indemnity,
+ or other liability obligations and/or rights consistent with this
+ License. However, in accepting such obligations, You may act only
+ on Your own behalf and on Your sole responsibility, not on behalf
+ of any other Contributor, and only if You agree to indemnify,
+ defend, and hold each Contributor harmless for any liability
+ incurred by, or claims asserted against, such Contributor by reason
+ of your accepting any such warranty or additional liability.
+
+ END OF TERMS AND CONDITIONS
+
+ APPENDIX: How to apply the Apache License to your work.
+
+ To apply the Apache License to your work, attach the following
+ boilerplate notice, with the fields enclosed by brackets "[]"
+ replaced with your own identifying information. (Don't include
+ the brackets!) The text should be enclosed in the appropriate
+ comment syntax for the file format. We also recommend that a
+ file or class name and description of purpose be included on the
+ same "printed page" as the copyright notice for easier
+ identification within third-party archives.
+
+ Copyright [yyyy] [name of copyright owner]
+
+ Licensed 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.
diff --git a/lib/netty-all-4.1.28.Final.jar b/lib/netty-all-4.1.37.Final.jar
similarity index 56%
rename from lib/netty-all-4.1.28.Final.jar
rename to lib/netty-all-4.1.37.Final.jar
index 058662ecc9..93cff04768 100644
Binary files a/lib/netty-all-4.1.28.Final.jar and b/lib/netty-all-4.1.37.Final.jar differ
diff --git a/lib/netty-tcnative-boringssl-static-2.0.25.Final.jar b/lib/netty-tcnative-boringssl-static-2.0.25.Final.jar
new file mode 100644
index 0000000000..954627fb73
Binary files /dev/null and b/lib/netty-tcnative-boringssl-static-2.0.25.Final.jar differ
diff --git a/src/java/org/apache/cassandra/net/InboundConnectionInitiator.java b/src/java/org/apache/cassandra/net/InboundConnectionInitiator.java
index d26abfdb25..c390ba4287 100644
--- a/src/java/org/apache/cassandra/net/InboundConnectionInitiator.java
+++ b/src/java/org/apache/cassandra/net/InboundConnectionInitiator.java
@@ -132,6 +132,8 @@ public class InboundConnectionInitiator
ServerBootstrap bootstrap = initializer.settings.socketFactory
.newServerBootstrap()
.option(ChannelOption.SO_BACKLOG, 1 << 9)
+ .option(ChannelOption.ALLOCATOR, GlobalBufferPoolAllocator.instance)
+ .option(ChannelOption.SO_REUSEADDR, true)
.childHandler(initializer);
int socketReceiveBufferSizeInBytes = initializer.settings.socketReceiveBufferSizeInBytes;
diff --git a/src/java/org/apache/cassandra/net/OutboundConnectionInitiator.java b/src/java/org/apache/cassandra/net/OutboundConnectionInitiator.java
index a63ccf9285..fdfb2dfa74 100644
--- a/src/java/org/apache/cassandra/net/OutboundConnectionInitiator.java
+++ b/src/java/org/apache/cassandra/net/OutboundConnectionInitiator.java
@@ -171,13 +171,15 @@ public class OutboundConnectionInitiator 0)
bootstrap.option(ChannelOption.SO_SNDBUF, settings.socketSendBufferSizeInBytes);
diff --git a/src/java/org/apache/cassandra/net/SocketFactory.java b/src/java/org/apache/cassandra/net/SocketFactory.java
index 18bb0d5c70..062c44ba63 100644
--- a/src/java/org/apache/cassandra/net/SocketFactory.java
+++ b/src/java/org/apache/cassandra/net/SocketFactory.java
@@ -18,13 +18,14 @@
package org.apache.cassandra.net;
import java.io.IOException;
-import java.lang.reflect.Field;
import java.net.ConnectException;
import java.net.InetSocketAddress;
import java.nio.channels.ClosedChannelException;
+import java.nio.channels.spi.SelectorProvider;
import java.util.List;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
+import java.util.concurrent.ThreadFactory;
import java.util.concurrent.TimeoutException;
import javax.annotation.Nullable;
import javax.net.ssl.SSLEngine;
@@ -37,11 +38,11 @@ import org.slf4j.LoggerFactory;
import io.netty.bootstrap.Bootstrap;
import io.netty.bootstrap.ServerBootstrap;
import io.netty.channel.Channel;
-import io.netty.channel.ChannelOption;
+import io.netty.channel.ChannelFactory;
+import io.netty.channel.DefaultSelectStrategyFactory;
import io.netty.channel.EventLoop;
import io.netty.channel.EventLoopGroup;
-import io.netty.channel.MultithreadEventLoopGroup;
-import io.netty.channel.SingleThreadEventLoop;
+import io.netty.channel.ServerChannel;
import io.netty.channel.epoll.EpollChannelOption;
import io.netty.channel.epoll.EpollEventLoopGroup;
import io.netty.channel.epoll.EpollServerSocketChannel;
@@ -53,10 +54,10 @@ import io.netty.channel.unix.Errors;
import io.netty.handler.ssl.OpenSsl;
import io.netty.handler.ssl.SslContext;
import io.netty.handler.ssl.SslHandler;
+import io.netty.util.concurrent.DefaultEventExecutorChooserFactory;
import io.netty.util.concurrent.DefaultThreadFactory;
-import io.netty.util.concurrent.EventExecutor;
-import io.netty.util.concurrent.MultithreadEventExecutorGroup;
-import io.netty.util.concurrent.SingleThreadEventExecutor;
+import io.netty.util.concurrent.RejectedExecutionHandlers;
+import io.netty.util.concurrent.ThreadPerTaskExecutor;
import io.netty.util.internal.logging.InternalLoggerFactory;
import io.netty.util.internal.logging.Slf4JLoggerFactory;
import org.apache.cassandra.concurrent.NamedThreadFactory;
@@ -82,8 +83,85 @@ public final class SocketFactory
private static final int EVENT_THREADS = Integer.getInteger(Config.PROPERTY_PREFIX + "internode-event-threads", FBUtilities.getAvailableProcessors());
- public enum Provider { EPOLL, NIO }
- private static final Provider DEFAULT_PROVIDER = NativeTransportService.useEpoll() ? Provider.EPOLL : Provider.NIO;
+ /**
+ * The default task queue used by {@code NioEventLoop} and {@code EpollEventLoop} is {@code MpscUnboundedArrayQueue},
+ * provided by JCTools. While efficient, it has an undesirable quality for a queue backing an event loop: it is
+ * not non-blocking, and can cause the event loop to busy-spin while waiting for a partially completed task
+ * offer, if the producer thread has been suspended mid-offer.
+ *
+ * As it happens, however, we have an MPSC queue implementation that is perfectly fit for this purpose -
+ * {@link ManyToOneConcurrentLinkedQueue}, that is non-blocking, and already used throughout the codebase,
+ * that we can and do use here as well.
+ */
+ enum Provider
+ {
+ NIO
+ {
+ @Override
+ NioEventLoopGroup makeEventLoopGroup(int threadCount, ThreadFactory threadFactory)
+ {
+ return new NioEventLoopGroup(threadCount,
+ new ThreadPerTaskExecutor(threadFactory),
+ DefaultEventExecutorChooserFactory.INSTANCE,
+ SelectorProvider.provider(),
+ DefaultSelectStrategyFactory.INSTANCE,
+ RejectedExecutionHandlers.reject(),
+ capacity -> new ManyToOneConcurrentLinkedQueue<>());
+ }
+
+ @Override
+ ChannelFactory clientChannelFactory()
+ {
+ return NioSocketChannel::new;
+ }
+
+ @Override
+ ChannelFactory serverChannelFactory()
+ {
+ return NioServerSocketChannel::new;
+ }
+ },
+ EPOLL
+ {
+ @Override
+ EpollEventLoopGroup makeEventLoopGroup(int threadCount, ThreadFactory threadFactory)
+ {
+ return new EpollEventLoopGroup(threadCount,
+ new ThreadPerTaskExecutor(threadFactory),
+ DefaultEventExecutorChooserFactory.INSTANCE,
+ DefaultSelectStrategyFactory.INSTANCE,
+ RejectedExecutionHandlers.reject(),
+ capacity -> new ManyToOneConcurrentLinkedQueue<>());
+ }
+
+ @Override
+ ChannelFactory clientChannelFactory()
+ {
+ return EpollSocketChannel::new;
+ }
+
+ @Override
+ ChannelFactory serverChannelFactory()
+ {
+ return EpollServerSocketChannel::new;
+ }
+ };
+
+ EventLoopGroup makeEventLoopGroup(int threadCount, String threadNamePrefix)
+ {
+ logger.debug("using netty {} event loop for pool prefix {}", name(), threadNamePrefix);
+ return makeEventLoopGroup(threadCount, new DefaultThreadFactory(threadNamePrefix, true));
+ }
+
+ abstract EventLoopGroup makeEventLoopGroup(int threadCount, ThreadFactory threadFactory);
+ abstract ChannelFactory extends Channel> clientChannelFactory();
+ abstract ChannelFactory extends ServerChannel> serverChannelFactory();
+
+ static Provider optimalProvider()
+ {
+ return NativeTransportService.useEpoll() ? EPOLL : NIO;
+ }
+ }
/** a useful addition for debugging; simply set to true to get more data in your logs */
static final boolean WIRETRACE = false;
@@ -93,94 +171,42 @@ public final class SocketFactory
InternalLoggerFactory.setDefaultFactory(Slf4JLoggerFactory.INSTANCE);
}
+ private final Provider provider;
private final EventLoopGroup acceptGroup;
private final EventLoopGroup defaultGroup;
// we need a separate EventLoopGroup for outbound streaming because sendFile is blocking
private final EventLoopGroup outboundStreamingGroup;
final ExecutorService synchronousWorkExecutor = Executors.newCachedThreadPool(new NamedThreadFactory("Messaging-SynchronousWork"));
- SocketFactory() { this(DEFAULT_PROVIDER); }
+ SocketFactory()
+ {
+ this(Provider.optimalProvider());
+ }
+
SocketFactory(Provider provider)
{
- this.acceptGroup = getEventLoopGroup(provider, 1, "Messaging-AcceptLoop");
- this.defaultGroup = getEventLoopGroup(provider, EVENT_THREADS, NamedThreadFactory.globalPrefix() + "Messaging-EventLoop");
- this.outboundStreamingGroup = getEventLoopGroup(provider, EVENT_THREADS, "Streaming-EventLoop");
- assert provider == providerOf(acceptGroup)
- && provider == providerOf(defaultGroup)
- && provider == providerOf(outboundStreamingGroup);
+ this.provider = provider;
+ this.acceptGroup = provider.makeEventLoopGroup(1, "Messaging-AcceptLoop");
+ this.defaultGroup = provider.makeEventLoopGroup(EVENT_THREADS, NamedThreadFactory.globalPrefix() + "Messaging-EventLoop");
+ this.outboundStreamingGroup = provider.makeEventLoopGroup(EVENT_THREADS, "Streaming-EventLoop");
}
- private static EventLoopGroup getEventLoopGroup(Provider provider, int threadCount, String threadNamePrefix)
- {
- switch (provider)
- {
- case EPOLL:
- logger.debug("using netty epoll event loop for pool prefix {}", threadNamePrefix);
- return overwriteMPSCQueues(new EpollEventLoopGroup(threadCount, new DefaultThreadFactory(threadNamePrefix, true)));
- case NIO:
- logger.debug("using netty nio event loop for pool prefix {}", threadNamePrefix);
- return overwriteMPSCQueues(new NioEventLoopGroup(threadCount, new DefaultThreadFactory(threadNamePrefix, true)));
- default:
- throw new IllegalStateException();
- }
- }
-
- private static Provider providerOf(EventLoopGroup eventLoopGroup)
- {
- while (eventLoopGroup instanceof SingleThreadEventLoop)
- eventLoopGroup = ((SingleThreadEventLoop) eventLoopGroup).parent();
-
- if (eventLoopGroup instanceof EpollEventLoopGroup)
- return Provider.EPOLL;
- if (eventLoopGroup instanceof NioEventLoopGroup)
- return Provider.NIO;
- throw new IllegalStateException();
- }
-
- static Bootstrap newBootstrap(EventLoop eventLoop, int tcpUserTimeoutInMS)
+ Bootstrap newClientBootstrap(EventLoop eventLoop, int tcpUserTimeoutInMS)
{
if (eventLoop == null)
throw new IllegalArgumentException("must provide eventLoop");
- Bootstrap bootstrap = new Bootstrap()
- .group(eventLoop)
- .option(ChannelOption.ALLOCATOR, GlobalBufferPoolAllocator.instance)
- .option(ChannelOption.SO_KEEPALIVE, true);
+ Bootstrap bootstrap = new Bootstrap().group(eventLoop).channelFactory(provider.clientChannelFactory());
+
+ if (provider == Provider.EPOLL)
+ bootstrap.option(EpollChannelOption.TCP_USER_TIMEOUT, tcpUserTimeoutInMS);
- switch (providerOf(eventLoop))
- {
- case EPOLL:
- bootstrap.channel(EpollSocketChannel.class);
- bootstrap.option(EpollChannelOption.TCP_USER_TIMEOUT, tcpUserTimeoutInMS);
- break;
- case NIO:
- bootstrap.channel(NioSocketChannel.class);
- }
return bootstrap;
}
ServerBootstrap newServerBootstrap()
{
- return newServerBootstrap(acceptGroup, defaultGroup);
- }
-
- private static ServerBootstrap newServerBootstrap(EventLoopGroup acceptGroup, EventLoopGroup defaultGroup)
- {
- ServerBootstrap bootstrap = new ServerBootstrap()
- .group(acceptGroup, defaultGroup)
- .option(ChannelOption.ALLOCATOR, GlobalBufferPoolAllocator.instance)
- .option(ChannelOption.SO_REUSEADDR, true);
-
- switch (providerOf(defaultGroup))
- {
- case EPOLL:
- bootstrap.channel(EpollServerSocketChannel.class);
- break;
- case NIO:
- bootstrap.channel(NioServerSocketChannel.class);
- }
-
- return bootstrap;
+ return new ServerBootstrap().group(acceptGroup, defaultGroup).channelFactory(provider.serverChannelFactory());
}
/**
@@ -270,58 +296,4 @@ public final class SocketFactory
{
return from + "->" + to + '-' + type + '-' + id;
}
-
- /**
- * The default task queue used by {@code NioEventLoop} and {@code EpollEventLoop} is {@code MpscUnboundedArrayQueue},
- * provided by JCTools. While efficient, it has an undesirable quality for a queue backing an event loop: it is
- * not non-blocking, and can cause the event loop to busy-spin while waiting for a partially completed task
- * offer, if the producer thread has been suspended mid-offer. Sadly, there is currently no way to work around
- * this behaviour in application-logic.
- *
- * As it happens, however, we have an MPSC queue implementation that is perfectly fit for this purpose -
- * {@link ManyToOneConcurrentLinkedQueue}, that is non-blocking, and already used throughout the codebase.
- *
- * Unfortunately, there is no Netty API or to override the default queue, so we have to resort to reflection,
- * for now.
- *
- * We filed a Netty issue asking for this capability to be provided cleanly:
- * https://github.com/netty/netty/issues/9105, and hopefully Netty will implement it some day. When and if
- * that happens, this reflection-based workaround should be removed.
- */
- private static EventLoopGroup overwriteMPSCQueues(MultithreadEventLoopGroup eventLoopGroup)
- {
- try
- {
- for (EventExecutor eventExecutor : (EventExecutor[]) childrenField.get(eventLoopGroup))
- {
- SingleThreadEventLoop eventLoop = (SingleThreadEventLoop) eventExecutor;
- taskQueueField.set(eventLoop, new ManyToOneConcurrentLinkedQueue<>());
- tailTasksField.set(eventLoop, new ManyToOneConcurrentLinkedQueue<>());
- }
- return eventLoopGroup;
- }
- catch (IllegalAccessException e)
- {
- throw new IllegalStateException(e);
- }
- }
-
- private static final Field childrenField, taskQueueField, tailTasksField;
- static
- {
- try
- {
- childrenField = MultithreadEventExecutorGroup.class.getDeclaredField("children");
- taskQueueField = SingleThreadEventExecutor.class.getDeclaredField("taskQueue");
- tailTasksField = SingleThreadEventLoop.class.getDeclaredField("tailTasks");
-
- childrenField.setAccessible(true);
- taskQueueField.setAccessible(true);
- tailTasksField.setAccessible(true);
- }
- catch (NoSuchFieldException e)
- {
- throw new IllegalStateException(e);
- }
- }
}
diff --git a/src/java/org/apache/cassandra/security/SSLFactory.java b/src/java/org/apache/cassandra/security/SSLFactory.java
index 667496388c..2ccb1268e6 100644
--- a/src/java/org/apache/cassandra/security/SSLFactory.java
+++ b/src/java/org/apache/cassandra/security/SSLFactory.java
@@ -240,10 +240,12 @@ public final class SSLFactory
* Get a netty {@link SslContext} instance.
*/
@VisibleForTesting
- static SslContext getOrCreateSslContext(EncryptionOptions options, boolean buildTruststore,
- SocketType socketType, boolean useOpenSsl) throws IOException
+ static SslContext getOrCreateSslContext(EncryptionOptions options,
+ boolean buildTruststore,
+ SocketType socketType,
+ boolean useOpenSsl) throws IOException
{
- CacheKey key = new CacheKey(options, socketType);
+ CacheKey key = new CacheKey(options, socketType, useOpenSsl);
SslContext sslContext;
sslContext = cachedSslContexts.get(key);
@@ -413,11 +415,13 @@ public final class SSLFactory
{
private final EncryptionOptions encryptionOptions;
private final SocketType socketType;
+ private final boolean useOpenSSL;
- public CacheKey(EncryptionOptions encryptionOptions, SocketType socketType)
+ public CacheKey(EncryptionOptions encryptionOptions, SocketType socketType, boolean useOpenSSL)
{
this.encryptionOptions = encryptionOptions;
this.socketType = socketType;
+ this.useOpenSSL = useOpenSSL;
}
public boolean equals(Object o)
@@ -426,6 +430,7 @@ public final class SSLFactory
if (o == null || getClass() != o.getClass()) return false;
CacheKey cacheKey = (CacheKey) o;
return (socketType == cacheKey.socketType &&
+ useOpenSSL == cacheKey.useOpenSSL &&
Objects.equals(encryptionOptions, cacheKey.encryptionOptions));
}
@@ -434,6 +439,7 @@ public final class SSLFactory
int result = 0;
result += 31 * socketType.hashCode();
result += 31 * encryptionOptions.hashCode();
+ result += 31 * Boolean.hashCode(useOpenSSL);
return result;
}
}
diff --git a/src/java/org/apache/cassandra/service/NativeTransportService.java b/src/java/org/apache/cassandra/service/NativeTransportService.java
index 79acab1b02..79caafcc7f 100644
--- a/src/java/org/apache/cassandra/service/NativeTransportService.java
+++ b/src/java/org/apache/cassandra/service/NativeTransportService.java
@@ -153,7 +153,7 @@ public class NativeTransportService
final boolean enableEpoll = Boolean.parseBoolean(System.getProperty("cassandra.native.epoll.enabled", "true"));
if (enableEpoll && !Epoll.isAvailable() && NativeLibrary.osType == NativeLibrary.OSType.LINUX)
- logger.warn("epoll not available {}", Epoll.unavailabilityCause());
+ logger.warn("epoll not available", Epoll.unavailabilityCause());
return enableEpoll && Epoll.isAvailable();
}