mirror of https://github.com/apache/cassandra
Merge branch 'cassandra-3.11' into cassandra-4.0
This commit is contained in:
commit
57293e8281
|
|
@ -56,6 +56,8 @@
|
||||||
<resolver:remoterepos id="all">
|
<resolver:remoterepos id="all">
|
||||||
<remoterepo id="resolver-central" url="${artifact.remoteRepository.central}"/>
|
<remoterepo id="resolver-central" url="${artifact.remoteRepository.central}"/>
|
||||||
<remoterepo id="resolver-apache" url="${artifact.remoteRepository.apache}"/>
|
<remoterepo id="resolver-apache" url="${artifact.remoteRepository.apache}"/>
|
||||||
|
<!-- Snapshots are not allowed, but for feature branches they may be needed, so uncomment the below to allow snapshots to work -->
|
||||||
|
<!-- <remoterepo id="resolver-apache-snapshot" url="https://repository.apache.org/content/repositories/snapshots" releases="false" snapshots="true" updates="always" checksums="fail" /> -->
|
||||||
</resolver:remoterepos>
|
</resolver:remoterepos>
|
||||||
|
|
||||||
<macrodef name="resolve">
|
<macrodef name="resolve">
|
||||||
|
|
|
||||||
|
|
@ -206,6 +206,7 @@
|
||||||
<string>--add-exports java.management.rmi/com.sun.jmx.remote.internal.rmi=ALL-UNNAMED</string>
|
<string>--add-exports java.management.rmi/com.sun.jmx.remote.internal.rmi=ALL-UNNAMED</string>
|
||||||
<string>--add-exports java.rmi/sun.rmi.registry=ALL-UNNAMED</string>
|
<string>--add-exports java.rmi/sun.rmi.registry=ALL-UNNAMED</string>
|
||||||
<string>--add-exports java.rmi/sun.rmi.server=ALL-UNNAMED</string>
|
<string>--add-exports java.rmi/sun.rmi.server=ALL-UNNAMED</string>
|
||||||
|
<string>--add-exports java.rmi/sun.rmi.transport.tcp=ALL-UNNAMED</string>
|
||||||
<string>--add-exports java.sql/java.sql=ALL-UNNAMED</string>
|
<string>--add-exports java.sql/java.sql=ALL-UNNAMED</string>
|
||||||
|
|
||||||
<string>--add-opens java.base/java.lang.module=ALL-UNNAMED</string>
|
<string>--add-opens java.base/java.lang.module=ALL-UNNAMED</string>
|
||||||
|
|
@ -217,6 +218,7 @@
|
||||||
<string>--add-opens java.base/jdk.internal.module=ALL-UNNAMED</string>
|
<string>--add-opens java.base/jdk.internal.module=ALL-UNNAMED</string>
|
||||||
<string>--add-opens java.base/jdk.internal.util.jar=ALL-UNNAMED</string>
|
<string>--add-opens java.base/jdk.internal.util.jar=ALL-UNNAMED</string>
|
||||||
<string>--add-opens jdk.management/com.sun.management.internal=ALL-UNNAMED</string>
|
<string>--add-opens jdk.management/com.sun.management.internal=ALL-UNNAMED</string>
|
||||||
|
|
||||||
</resources>
|
</resources>
|
||||||
<pathconvert property="_jvm_args_concat" refid="_jvm11_arg_items" pathsep=" "/>
|
<pathconvert property="_jvm_args_concat" refid="_jvm11_arg_items" pathsep=" "/>
|
||||||
<condition property="java11-jvmargs" value="${_jvm_args_concat}" else="">
|
<condition property="java11-jvmargs" value="${_jvm_args_concat}" else="">
|
||||||
|
|
@ -257,7 +259,7 @@
|
||||||
</condition>
|
</condition>
|
||||||
|
|
||||||
<!-- needed to compile org.apache.cassandra.utils.JMXServerUtils -->
|
<!-- needed to compile org.apache.cassandra.utils.JMXServerUtils -->
|
||||||
<condition property="jdk11-javac-exports" value="--add-exports java.rmi/sun.rmi.registry=ALL-UNNAMED" else="">
|
<condition property="jdk11-javac-exports" value="--add-exports java.rmi/sun.rmi.registry=ALL-UNNAMED --add-exports java.rmi/sun.rmi.transport.tcp=ALL-UNNAMED" else="">
|
||||||
<not>
|
<not>
|
||||||
<equals arg1="${ant.java.version}" arg2="1.8"/>
|
<equals arg1="${ant.java.version}" arg2="1.8"/>
|
||||||
</not>
|
</not>
|
||||||
|
|
@ -554,7 +556,7 @@
|
||||||
<dependency groupId="com.google.code.java-allocation-instrumenter" artifactId="java-allocation-instrumenter" version="${allocation-instrumenter.version}" scope="test">
|
<dependency groupId="com.google.code.java-allocation-instrumenter" artifactId="java-allocation-instrumenter" version="${allocation-instrumenter.version}" scope="test">
|
||||||
<exclusion groupId="com.google.guava" artifactId="guava"/>
|
<exclusion groupId="com.google.guava" artifactId="guava"/>
|
||||||
</dependency>
|
</dependency>
|
||||||
<dependency groupId="org.apache.cassandra" artifactId="dtest-api" version="0.0.13" scope="test"/>
|
<dependency groupId="org.apache.cassandra" artifactId="dtest-api" version="0.0.15" scope="test"/>
|
||||||
<dependency groupId="org.reflections" artifactId="reflections" version="0.10.2" scope="test"/>
|
<dependency groupId="org.reflections" artifactId="reflections" version="0.10.2" scope="test"/>
|
||||||
<dependency groupId="org.apache.hadoop" artifactId="hadoop-core" version="1.0.3" scope="provided">
|
<dependency groupId="org.apache.hadoop" artifactId="hadoop-core" version="1.0.3" scope="provided">
|
||||||
<exclusion groupId="org.mortbay.jetty" artifactId="servlet-api"/>
|
<exclusion groupId="org.mortbay.jetty" artifactId="servlet-api"/>
|
||||||
|
|
@ -1303,6 +1305,7 @@
|
||||||
<pathelement location="${fqltool.build.classes}"/>
|
<pathelement location="${fqltool.build.classes}"/>
|
||||||
</classpath>
|
</classpath>
|
||||||
<compilerarg value="-XDignore.symbol.file"/>
|
<compilerarg value="-XDignore.symbol.file"/>
|
||||||
|
<compilerarg line="${jdk11-javac-exports}"/>
|
||||||
<src path="${test.anttasks.src}"/>
|
<src path="${test.anttasks.src}"/>
|
||||||
<src path="${test.unit.src}"/>
|
<src path="${test.unit.src}"/>
|
||||||
<src path="${test.long.src}"/>
|
<src path="${test.long.src}"/>
|
||||||
|
|
@ -2010,7 +2013,7 @@
|
||||||
<echo file=".idea/compiler.xml"><![CDATA[<?xml version="1.0" encoding="UTF-8"?>
|
<echo file=".idea/compiler.xml"><![CDATA[<?xml version="1.0" encoding="UTF-8"?>
|
||||||
<project version="4">
|
<project version="4">
|
||||||
<component name="JavacSettings">
|
<component name="JavacSettings">
|
||||||
<option name="ADDITIONAL_OPTIONS_STRING" value="--add-exports java.rmi/sun.rmi.registry=ALL-UNNAMED" />
|
<option name="ADDITIONAL_OPTIONS_STRING" value="--add-exports java.rmi/sun.rmi.registry=ALL-UNNAMED --add-exports java.rmi/sun.rmi.transport.tcp=ALL-UNNAMED" />
|
||||||
</component>
|
</component>
|
||||||
</project>]]></echo>
|
</project>]]></echo>
|
||||||
</target>
|
</target>
|
||||||
|
|
|
||||||
|
|
@ -67,7 +67,7 @@
|
||||||
<plugin>
|
<plugin>
|
||||||
<groupId>org.apache.maven.plugins</groupId>
|
<groupId>org.apache.maven.plugins</groupId>
|
||||||
<artifactId>maven-shade-plugin</artifactId>
|
<artifactId>maven-shade-plugin</artifactId>
|
||||||
<version>3.2.1</version>
|
<version>3.4.1</version>
|
||||||
|
|
||||||
<configuration>
|
<configuration>
|
||||||
<createSourcesJar>false</createSourcesJar>
|
<createSourcesJar>false</createSourcesJar>
|
||||||
|
|
|
||||||
|
|
@ -87,6 +87,19 @@ public enum CassandraRelevantProperties
|
||||||
*/
|
*/
|
||||||
COM_SUN_MANAGEMENT_JMXREMOTE_AUTHENTICATE ("com.sun.management.jmxremote.authenticate"),
|
COM_SUN_MANAGEMENT_JMXREMOTE_AUTHENTICATE ("com.sun.management.jmxremote.authenticate"),
|
||||||
|
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Controls the JMX server threadpool keap-alive time.
|
||||||
|
* Should only be set by in-jvm dtests.
|
||||||
|
*/
|
||||||
|
SUN_RMI_TRANSPORT_TCP_THREADKEEPALIVETIME("sun.rmi.transport.tcp.threadKeepAliveTime"),
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Controls the distributed garbage collector lease time for JMX objects.
|
||||||
|
* Should only be set by in-jvm dtests.
|
||||||
|
*/
|
||||||
|
JAVA_RMI_DGC_LEASE_VALUE_IN_JVM_DTEST("java.rmi.dgc.leaseValue"),
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* The port number to which the RMI connector will be bound - com.sun.management.jmxremote.rmi.port.
|
* The port number to which the RMI connector will be bound - com.sun.management.jmxremote.rmi.port.
|
||||||
* An Integer object that represents the value of the second argument is returned
|
* An Integer object that represents the value of the second argument is returned
|
||||||
|
|
@ -196,6 +209,9 @@ public enum CassandraRelevantProperties
|
||||||
/** what class to use for mbean registeration */
|
/** what class to use for mbean registeration */
|
||||||
MBEAN_REGISTRATION_CLASS("org.apache.cassandra.mbean_registration_class"),
|
MBEAN_REGISTRATION_CLASS("org.apache.cassandra.mbean_registration_class"),
|
||||||
|
|
||||||
|
/** This property indicates if the code is running under the in-jvm dtest framework */
|
||||||
|
DTEST_IS_IN_JVM_DTEST("org.apache.cassandra.dtest.is_in_jvm_dtest"),
|
||||||
|
|
||||||
MIGRATION_DELAY("cassandra.migration_delay_ms", "60000"),
|
MIGRATION_DELAY("cassandra.migration_delay_ms", "60000"),
|
||||||
/** Defines how often schema definitions are pulled from the other nodes */
|
/** Defines how often schema definitions are pulled from the other nodes */
|
||||||
SCHEMA_PULL_INTERVAL_MS("cassandra.schema_pull_interval_ms", "60000"),
|
SCHEMA_PULL_INTERVAL_MS("cassandra.schema_pull_interval_ms", "60000"),
|
||||||
|
|
|
||||||
|
|
@ -146,18 +146,16 @@ public class GCInspector implements NotificationListener, GCInspectorMXBean
|
||||||
|
|
||||||
public GCInspector()
|
public GCInspector()
|
||||||
{
|
{
|
||||||
MBeanServer mbs = ManagementFactory.getPlatformMBeanServer();
|
|
||||||
|
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
ObjectName gcName = new ObjectName(ManagementFactory.GARBAGE_COLLECTOR_MXBEAN_DOMAIN_TYPE + ",*");
|
ObjectName gcName = new ObjectName(ManagementFactory.GARBAGE_COLLECTOR_MXBEAN_DOMAIN_TYPE + ",*");
|
||||||
for (ObjectName name : mbs.queryNames(gcName, null))
|
for (ObjectName name : MBeanWrapper.instance.queryNames(gcName, null))
|
||||||
{
|
{
|
||||||
GarbageCollectorMXBean gc = ManagementFactory.newPlatformMXBeanProxy(mbs, name.getCanonicalName(), GarbageCollectorMXBean.class);
|
GarbageCollectorMXBean gc = ManagementFactory.newPlatformMXBeanProxy(MBeanWrapper.instance.getMBeanServer(), name.getCanonicalName(), GarbageCollectorMXBean.class);
|
||||||
gcStates.put(gc.getName(), new GCState(gc, assumeGCIsPartiallyConcurrent(gc), assumeGCIsOldGen(gc)));
|
gcStates.put(gc.getName(), new GCState(gc, assumeGCIsPartiallyConcurrent(gc), assumeGCIsOldGen(gc)));
|
||||||
}
|
}
|
||||||
ObjectName me = new ObjectName(MBEAN_NAME);
|
ObjectName me = new ObjectName(MBEAN_NAME);
|
||||||
if (!mbs.isRegistered(me))
|
if (!MBeanWrapper.instance.isRegistered(me))
|
||||||
MBeanWrapper.instance.registerMBean(this, new ObjectName(MBEAN_NAME));
|
MBeanWrapper.instance.registerMBean(this, new ObjectName(MBEAN_NAME));
|
||||||
}
|
}
|
||||||
catch (MalformedObjectNameException | IOException e)
|
catch (MalformedObjectNameException | IOException e)
|
||||||
|
|
|
||||||
|
|
@ -16,13 +16,15 @@
|
||||||
* limitations under the License.
|
* limitations under the License.
|
||||||
*/
|
*/
|
||||||
|
|
||||||
package org.apache.cassandra.distributed.shared;
|
package org.apache.cassandra.utils;
|
||||||
|
|
||||||
import java.lang.annotation.ElementType;
|
import java.lang.annotation.ElementType;
|
||||||
import java.lang.annotation.Retention;
|
import java.lang.annotation.Retention;
|
||||||
import java.lang.annotation.RetentionPolicy;
|
import java.lang.annotation.RetentionPolicy;
|
||||||
import java.lang.annotation.Target;
|
import java.lang.annotation.Target;
|
||||||
|
|
||||||
|
import org.apache.cassandra.utils.Shared;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Tells jvm-dtest that a class should be isolated and loaded into the instance class loader.
|
* Tells jvm-dtest that a class should be isolated and loaded into the instance class loader.
|
||||||
*
|
*
|
||||||
|
|
@ -29,12 +29,14 @@ import java.net.Inet6Address;
|
||||||
import java.net.InetAddress;
|
import java.net.InetAddress;
|
||||||
import java.rmi.AccessException;
|
import java.rmi.AccessException;
|
||||||
import java.rmi.AlreadyBoundException;
|
import java.rmi.AlreadyBoundException;
|
||||||
|
import java.rmi.NoSuchObjectException;
|
||||||
import java.rmi.NotBoundException;
|
import java.rmi.NotBoundException;
|
||||||
import java.rmi.Remote;
|
import java.rmi.Remote;
|
||||||
import java.rmi.RemoteException;
|
import java.rmi.RemoteException;
|
||||||
import java.rmi.registry.Registry;
|
import java.rmi.registry.Registry;
|
||||||
import java.rmi.server.RMIClientSocketFactory;
|
import java.rmi.server.RMIClientSocketFactory;
|
||||||
import java.rmi.server.RMIServerSocketFactory;
|
import java.rmi.server.RMIServerSocketFactory;
|
||||||
|
import java.rmi.server.UnicastRemoteObject;
|
||||||
import java.util.Arrays;
|
import java.util.Arrays;
|
||||||
import java.util.HashMap;
|
import java.util.HashMap;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
|
|
@ -252,7 +254,8 @@ public class JMXServerUtils
|
||||||
return env;
|
return env;
|
||||||
}
|
}
|
||||||
|
|
||||||
private static void logJmxServiceUrl(InetAddress serverAddress, int port)
|
@VisibleForTesting
|
||||||
|
public static void logJmxServiceUrl(InetAddress serverAddress, int port)
|
||||||
{
|
{
|
||||||
String urlTemplate = "service:jmx:rmi://%1$s/jndi/rmi://%1$s:%2$d/jmxrmi";
|
String urlTemplate = "service:jmx:rmi://%1$s/jndi/rmi://%1$s:%2$d/jmxrmi";
|
||||||
String hostName;
|
String hostName;
|
||||||
|
|
@ -325,11 +328,11 @@ public class JMXServerUtils
|
||||||
* Better to use the internal API than re-invent the wheel.
|
* Better to use the internal API than re-invent the wheel.
|
||||||
*/
|
*/
|
||||||
@SuppressWarnings("restriction")
|
@SuppressWarnings("restriction")
|
||||||
private static class JmxRegistry extends sun.rmi.registry.RegistryImpl {
|
public static class JmxRegistry extends sun.rmi.registry.RegistryImpl {
|
||||||
private final String lookupName;
|
private final String lookupName;
|
||||||
private Remote remoteServerStub;
|
private Remote remoteServerStub;
|
||||||
|
|
||||||
JmxRegistry(final int port,
|
public JmxRegistry(final int port,
|
||||||
final RMIClientSocketFactory csf,
|
final RMIClientSocketFactory csf,
|
||||||
RMIServerSocketFactory ssf,
|
RMIServerSocketFactory ssf,
|
||||||
final String lookupName) throws RemoteException {
|
final String lookupName) throws RemoteException {
|
||||||
|
|
@ -362,5 +365,24 @@ public class JMXServerUtils
|
||||||
public void setRemoteServerStub(Remote remoteServerStub) {
|
public void setRemoteServerStub(Remote remoteServerStub) {
|
||||||
this.remoteServerStub = remoteServerStub;
|
this.remoteServerStub = remoteServerStub;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Closes the underlying JMX registry by unexporting this instance.
|
||||||
|
* There is no reason to do this except for in-jvm dtests where we need
|
||||||
|
* to stop the registry, so we can start with a clean slate for future cluster
|
||||||
|
* builds, and the superclass never expects to be shut down and therefore doesn't
|
||||||
|
* handle this edge case at all.
|
||||||
|
*/
|
||||||
|
@VisibleForTesting
|
||||||
|
public void close() {
|
||||||
|
try
|
||||||
|
{
|
||||||
|
UnicastRemoteObject.unexportObject(this, true);
|
||||||
|
}
|
||||||
|
catch (NoSuchObjectException ignored)
|
||||||
|
{
|
||||||
|
// Ignore if it's already unexported
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -19,15 +19,21 @@
|
||||||
package org.apache.cassandra.utils;
|
package org.apache.cassandra.utils;
|
||||||
|
|
||||||
import java.lang.management.ManagementFactory;
|
import java.lang.management.ManagementFactory;
|
||||||
|
import java.util.Collections;
|
||||||
|
import java.util.Set;
|
||||||
|
import java.util.UUID;
|
||||||
import java.util.function.Consumer;
|
import java.util.function.Consumer;
|
||||||
import javax.management.MBeanServer;
|
import javax.management.MBeanServer;
|
||||||
|
import javax.management.MBeanServerFactory;
|
||||||
import javax.management.MalformedObjectNameException;
|
import javax.management.MalformedObjectNameException;
|
||||||
import javax.management.ObjectName;
|
import javax.management.ObjectName;
|
||||||
|
import javax.management.QueryExp;
|
||||||
|
|
||||||
import org.slf4j.Logger;
|
import org.slf4j.Logger;
|
||||||
import org.slf4j.LoggerFactory;
|
import org.slf4j.LoggerFactory;
|
||||||
|
|
||||||
import static org.apache.cassandra.config.CassandraRelevantProperties.IS_DISABLED_MBEAN_REGISTRATION;
|
import static org.apache.cassandra.config.CassandraRelevantProperties.ORG_APACHE_CASSANDRA_DISABLE_MBEAN_REGISTRATION;
|
||||||
|
import static org.apache.cassandra.config.CassandraRelevantProperties.DTEST_IS_IN_JVM_DTEST;
|
||||||
import static org.apache.cassandra.config.CassandraRelevantProperties.MBEAN_REGISTRATION_CLASS;
|
import static org.apache.cassandra.config.CassandraRelevantProperties.MBEAN_REGISTRATION_CLASS;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|
@ -36,23 +42,46 @@ import static org.apache.cassandra.config.CassandraRelevantProperties.MBEAN_REGI
|
||||||
*/
|
*/
|
||||||
public interface MBeanWrapper
|
public interface MBeanWrapper
|
||||||
{
|
{
|
||||||
static final Logger logger = LoggerFactory.getLogger(MBeanWrapper.class);
|
Logger logger = LoggerFactory.getLogger(MBeanWrapper.class);
|
||||||
|
|
||||||
static final MBeanWrapper instance = create();
|
MBeanWrapper instance = create();
|
||||||
|
|
||||||
static MBeanWrapper create()
|
static MBeanWrapper create()
|
||||||
{
|
{
|
||||||
if (IS_DISABLED_MBEAN_REGISTRATION.getBoolean())
|
// If we're running in the in-jvm dtest environment, always use the delegating
|
||||||
|
// mbean wrapper even if we start off with no-op, so it can be switched later
|
||||||
|
if (DTEST_IS_IN_JVM_DTEST.getBoolean())
|
||||||
|
{
|
||||||
|
return new DelegatingMbeanWrapper(getMBeanWrapper());
|
||||||
|
}
|
||||||
|
|
||||||
|
return getMBeanWrapper();
|
||||||
|
}
|
||||||
|
|
||||||
|
static MBeanWrapper getMBeanWrapper()
|
||||||
|
{
|
||||||
|
if (ORG_APACHE_CASSANDRA_DISABLE_MBEAN_REGISTRATION.getBoolean())
|
||||||
|
{
|
||||||
return new NoOpMBeanWrapper();
|
return new NoOpMBeanWrapper();
|
||||||
|
}
|
||||||
|
|
||||||
String klass = MBEAN_REGISTRATION_CLASS.getString();
|
String klass = MBEAN_REGISTRATION_CLASS.getString();
|
||||||
if (klass == null)
|
if (klass == null)
|
||||||
return new PlatformMBeanWrapper();
|
{
|
||||||
|
if (DTEST_IS_IN_JVM_DTEST.getBoolean())
|
||||||
|
{
|
||||||
|
return new NoOpMBeanWrapper();
|
||||||
|
}
|
||||||
|
else
|
||||||
|
{
|
||||||
|
return new PlatformMBeanWrapper();
|
||||||
|
}
|
||||||
|
}
|
||||||
return FBUtilities.construct(klass, "mbean");
|
return FBUtilities.construct(klass, "mbean");
|
||||||
}
|
}
|
||||||
|
|
||||||
// Passing true for graceful will log exceptions instead of rethrowing them
|
// Passing true for graceful will log exceptions instead of rethrowing them
|
||||||
public void registerMBean(Object obj, ObjectName mbeanName, OnException onException);
|
void registerMBean(Object obj, ObjectName mbeanName, OnException onException);
|
||||||
default void registerMBean(Object obj, ObjectName mbeanName)
|
default void registerMBean(Object obj, ObjectName mbeanName)
|
||||||
{
|
{
|
||||||
registerMBean(obj, mbeanName, OnException.THROW);
|
registerMBean(obj, mbeanName, OnException.THROW);
|
||||||
|
|
@ -62,7 +91,9 @@ public interface MBeanWrapper
|
||||||
{
|
{
|
||||||
ObjectName name = create(mbeanName, onException);
|
ObjectName name = create(mbeanName, onException);
|
||||||
if (name == null)
|
if (name == null)
|
||||||
|
{
|
||||||
return;
|
return;
|
||||||
|
}
|
||||||
registerMBean(obj, name, onException);
|
registerMBean(obj, name, onException);
|
||||||
}
|
}
|
||||||
default void registerMBean(Object obj, String mbeanName)
|
default void registerMBean(Object obj, String mbeanName)
|
||||||
|
|
@ -70,7 +101,7 @@ public interface MBeanWrapper
|
||||||
registerMBean(obj, mbeanName, OnException.THROW);
|
registerMBean(obj, mbeanName, OnException.THROW);
|
||||||
}
|
}
|
||||||
|
|
||||||
public boolean isRegistered(ObjectName mbeanName, OnException onException);
|
boolean isRegistered(ObjectName mbeanName, OnException onException);
|
||||||
default boolean isRegistered(ObjectName mbeanName)
|
default boolean isRegistered(ObjectName mbeanName)
|
||||||
{
|
{
|
||||||
return isRegistered(mbeanName, OnException.THROW);
|
return isRegistered(mbeanName, OnException.THROW);
|
||||||
|
|
@ -80,7 +111,9 @@ public interface MBeanWrapper
|
||||||
{
|
{
|
||||||
ObjectName name = create(mbeanName, onException);
|
ObjectName name = create(mbeanName, onException);
|
||||||
if (name == null)
|
if (name == null)
|
||||||
|
{
|
||||||
return false;
|
return false;
|
||||||
|
}
|
||||||
return isRegistered(name, onException);
|
return isRegistered(name, onException);
|
||||||
}
|
}
|
||||||
default boolean isRegistered(String mbeanName)
|
default boolean isRegistered(String mbeanName)
|
||||||
|
|
@ -88,7 +121,7 @@ public interface MBeanWrapper
|
||||||
return isRegistered(mbeanName, OnException.THROW);
|
return isRegistered(mbeanName, OnException.THROW);
|
||||||
}
|
}
|
||||||
|
|
||||||
public void unregisterMBean(ObjectName mbeanName, OnException onException);
|
void unregisterMBean(ObjectName mbeanName, OnException onException);
|
||||||
default void unregisterMBean(ObjectName mbeanName)
|
default void unregisterMBean(ObjectName mbeanName)
|
||||||
{
|
{
|
||||||
unregisterMBean(mbeanName, OnException.THROW);
|
unregisterMBean(mbeanName, OnException.THROW);
|
||||||
|
|
@ -98,7 +131,9 @@ public interface MBeanWrapper
|
||||||
{
|
{
|
||||||
ObjectName name = create(mbeanName, onException);
|
ObjectName name = create(mbeanName, onException);
|
||||||
if (name == null)
|
if (name == null)
|
||||||
|
{
|
||||||
return;
|
return;
|
||||||
|
}
|
||||||
unregisterMBean(name, onException);
|
unregisterMBean(name, onException);
|
||||||
}
|
}
|
||||||
default void unregisterMBean(String mbeanName)
|
default void unregisterMBean(String mbeanName)
|
||||||
|
|
@ -119,7 +154,11 @@ public interface MBeanWrapper
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
static class NoOpMBeanWrapper implements MBeanWrapper
|
Set<ObjectName> queryNames(ObjectName name, QueryExp query);
|
||||||
|
|
||||||
|
MBeanServer getMBeanServer();
|
||||||
|
|
||||||
|
class NoOpMBeanWrapper implements MBeanWrapper
|
||||||
{
|
{
|
||||||
public void registerMBean(Object obj, ObjectName mbeanName, OnException onException) {}
|
public void registerMBean(Object obj, ObjectName mbeanName, OnException onException) {}
|
||||||
public void registerMBean(Object obj, String mbeanName, OnException onException) {}
|
public void registerMBean(Object obj, String mbeanName, OnException onException) {}
|
||||||
|
|
@ -127,9 +166,11 @@ public interface MBeanWrapper
|
||||||
public boolean isRegistered(String mbeanName, OnException onException) { return false; }
|
public boolean isRegistered(String mbeanName, OnException onException) { return false; }
|
||||||
public void unregisterMBean(ObjectName mbeanName, OnException onException) {}
|
public void unregisterMBean(ObjectName mbeanName, OnException onException) {}
|
||||||
public void unregisterMBean(String mbeanName, OnException onException) {}
|
public void unregisterMBean(String mbeanName, OnException onException) {}
|
||||||
|
public Set<ObjectName> queryNames(ObjectName name, QueryExp query) {return Collections.emptySet(); }
|
||||||
|
public MBeanServer getMBeanServer() { return null; }
|
||||||
}
|
}
|
||||||
|
|
||||||
static class PlatformMBeanWrapper implements MBeanWrapper
|
class PlatformMBeanWrapper implements MBeanWrapper
|
||||||
{
|
{
|
||||||
private final MBeanServer mbs = ManagementFactory.getPlatformMBeanServer();
|
private final MBeanServer mbs = ManagementFactory.getPlatformMBeanServer();
|
||||||
public void registerMBean(Object obj, ObjectName mbeanName, OnException onException)
|
public void registerMBean(Object obj, ObjectName mbeanName, OnException onException)
|
||||||
|
|
@ -168,9 +209,158 @@ public interface MBeanWrapper
|
||||||
onException.handler.accept(e);
|
onException.handler.accept(e);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public Set<ObjectName> queryNames(ObjectName name, QueryExp query)
|
||||||
|
{
|
||||||
|
return mbs.queryNames(name, query);
|
||||||
|
}
|
||||||
|
|
||||||
|
public MBeanServer getMBeanServer()
|
||||||
|
{
|
||||||
|
return mbs;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
public enum OnException
|
class InstanceMBeanWrapper implements MBeanWrapper
|
||||||
|
{
|
||||||
|
private MBeanServer mbs;
|
||||||
|
public final UUID id = UUID.randomUUID();
|
||||||
|
|
||||||
|
public InstanceMBeanWrapper(String hostname)
|
||||||
|
{
|
||||||
|
mbs = MBeanServerFactory.createMBeanServer(hostname + "-" + id);
|
||||||
|
}
|
||||||
|
|
||||||
|
public void registerMBean(Object obj, ObjectName mbeanName, OnException onException)
|
||||||
|
{
|
||||||
|
try
|
||||||
|
{
|
||||||
|
mbs.registerMBean(obj, mbeanName);
|
||||||
|
}
|
||||||
|
catch (Exception e)
|
||||||
|
{
|
||||||
|
onException.handler.accept(e);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
public boolean isRegistered(ObjectName mbeanName, OnException onException)
|
||||||
|
{
|
||||||
|
try
|
||||||
|
{
|
||||||
|
return mbs.isRegistered(mbeanName);
|
||||||
|
}
|
||||||
|
catch (Exception e)
|
||||||
|
{
|
||||||
|
onException.handler.accept(e);
|
||||||
|
}
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void unregisterMBean(ObjectName mbeanName, OnException onException)
|
||||||
|
{
|
||||||
|
try
|
||||||
|
{
|
||||||
|
mbs.unregisterMBean(mbeanName);
|
||||||
|
}
|
||||||
|
catch (Exception e)
|
||||||
|
{
|
||||||
|
onException.handler.accept(e);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
public Set<ObjectName> queryNames(ObjectName name, QueryExp query)
|
||||||
|
{
|
||||||
|
return mbs.queryNames(name, query);
|
||||||
|
}
|
||||||
|
|
||||||
|
public MBeanServer getMBeanServer()
|
||||||
|
{
|
||||||
|
return mbs;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void close() {
|
||||||
|
mbs.queryNames(null, null).forEach(name -> {
|
||||||
|
try {
|
||||||
|
if (!name.getCanonicalName().contains("MBeanServerDelegate"))
|
||||||
|
{
|
||||||
|
mbs.unregisterMBean(name);
|
||||||
|
}
|
||||||
|
} catch (Throwable e) {
|
||||||
|
logger.debug("Could not unregister mbean {}", name.getCanonicalName());
|
||||||
|
}
|
||||||
|
});
|
||||||
|
MBeanServerFactory.releaseMBeanServer(mbs);
|
||||||
|
mbs = null;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
class DelegatingMbeanWrapper implements MBeanWrapper
|
||||||
|
{
|
||||||
|
MBeanWrapper delegate;
|
||||||
|
|
||||||
|
public DelegatingMbeanWrapper(MBeanWrapper mBeanWrapper)
|
||||||
|
{
|
||||||
|
delegate = mBeanWrapper;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setDelegate(MBeanWrapper wrapper) {
|
||||||
|
delegate = wrapper;
|
||||||
|
}
|
||||||
|
|
||||||
|
public MBeanWrapper getDelegate()
|
||||||
|
{
|
||||||
|
return delegate;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void registerMBean(Object obj, ObjectName mbeanName, OnException onException)
|
||||||
|
{
|
||||||
|
try
|
||||||
|
{
|
||||||
|
delegate.registerMBean(obj, mbeanName);
|
||||||
|
}
|
||||||
|
catch (Exception e)
|
||||||
|
{
|
||||||
|
onException.handler.accept(e);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
public boolean isRegistered(ObjectName mbeanName, OnException onException)
|
||||||
|
{
|
||||||
|
try
|
||||||
|
{
|
||||||
|
return delegate.isRegistered(mbeanName);
|
||||||
|
}
|
||||||
|
catch (Exception e)
|
||||||
|
{
|
||||||
|
onException.handler.accept(e);
|
||||||
|
}
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void unregisterMBean(ObjectName mbeanName, OnException onException)
|
||||||
|
{
|
||||||
|
try
|
||||||
|
{
|
||||||
|
delegate.unregisterMBean(mbeanName);
|
||||||
|
}
|
||||||
|
catch (Exception e)
|
||||||
|
{
|
||||||
|
onException.handler.accept(e);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
public Set<ObjectName> queryNames(ObjectName name, QueryExp query)
|
||||||
|
{
|
||||||
|
return delegate.queryNames(name, query);
|
||||||
|
}
|
||||||
|
|
||||||
|
public MBeanServer getMBeanServer()
|
||||||
|
{
|
||||||
|
return delegate.getMBeanServer();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
enum OnException
|
||||||
{
|
{
|
||||||
THROW(e -> { throw new RuntimeException(e); }),
|
THROW(e -> { throw new RuntimeException(e); }),
|
||||||
LOG(e -> { logger.error("Error in MBean wrapper: ", e); }),
|
LOG(e -> { logger.error("Error in MBean wrapper: ", e); }),
|
||||||
|
|
|
||||||
|
|
@ -39,10 +39,12 @@ import static org.apache.cassandra.utils.NativeLibrary.OSType.MAC;
|
||||||
import static org.apache.cassandra.utils.NativeLibrary.OSType.WINDOWS;
|
import static org.apache.cassandra.utils.NativeLibrary.OSType.WINDOWS;
|
||||||
import static org.apache.cassandra.utils.NativeLibrary.OSType.AIX;
|
import static org.apache.cassandra.utils.NativeLibrary.OSType.AIX;
|
||||||
|
|
||||||
|
@Shared
|
||||||
public final class NativeLibrary
|
public final class NativeLibrary
|
||||||
{
|
{
|
||||||
private static final Logger logger = LoggerFactory.getLogger(NativeLibrary.class);
|
private static final Logger logger = LoggerFactory.getLogger(NativeLibrary.class);
|
||||||
|
|
||||||
|
@Shared
|
||||||
public enum OSType
|
public enum OSType
|
||||||
{
|
{
|
||||||
LINUX,
|
LINUX,
|
||||||
|
|
|
||||||
|
|
@ -42,6 +42,7 @@ import com.sun.jna.Pointer;
|
||||||
* @see org.apache.cassandra.utils.NativeLibraryWrapper
|
* @see org.apache.cassandra.utils.NativeLibraryWrapper
|
||||||
* @see NativeLibrary
|
* @see NativeLibrary
|
||||||
*/
|
*/
|
||||||
|
@Shared
|
||||||
public class NativeLibraryDarwin implements NativeLibraryWrapper
|
public class NativeLibraryDarwin implements NativeLibraryWrapper
|
||||||
{
|
{
|
||||||
private static final Logger logger = LoggerFactory.getLogger(NativeLibraryDarwin.class);
|
private static final Logger logger = LoggerFactory.getLogger(NativeLibraryDarwin.class);
|
||||||
|
|
|
||||||
|
|
@ -42,6 +42,7 @@ import com.sun.jna.Pointer;
|
||||||
* @see org.apache.cassandra.utils.NativeLibraryWrapper
|
* @see org.apache.cassandra.utils.NativeLibraryWrapper
|
||||||
* @see NativeLibrary
|
* @see NativeLibrary
|
||||||
*/
|
*/
|
||||||
|
@Shared
|
||||||
public class NativeLibraryLinux implements NativeLibraryWrapper
|
public class NativeLibraryLinux implements NativeLibraryWrapper
|
||||||
{
|
{
|
||||||
private static boolean available;
|
private static boolean available;
|
||||||
|
|
|
||||||
|
|
@ -35,6 +35,7 @@ import com.sun.jna.Pointer;
|
||||||
* @see org.apache.cassandra.utils.NativeLibraryWrapper
|
* @see org.apache.cassandra.utils.NativeLibraryWrapper
|
||||||
* @see NativeLibrary
|
* @see NativeLibrary
|
||||||
*/
|
*/
|
||||||
|
@Shared
|
||||||
public class NativeLibraryWindows implements NativeLibraryWrapper
|
public class NativeLibraryWindows implements NativeLibraryWrapper
|
||||||
{
|
{
|
||||||
private static final Logger logger = LoggerFactory.getLogger(NativeLibraryWindows.class);
|
private static final Logger logger = LoggerFactory.getLogger(NativeLibraryWindows.class);
|
||||||
|
|
|
||||||
|
|
@ -24,6 +24,7 @@ import com.sun.jna.Pointer;
|
||||||
* An interface to implement for using OS specific native methods.
|
* An interface to implement for using OS specific native methods.
|
||||||
* @see NativeLibrary
|
* @see NativeLibrary
|
||||||
*/
|
*/
|
||||||
|
@Shared
|
||||||
interface NativeLibraryWrapper
|
interface NativeLibraryWrapper
|
||||||
{
|
{
|
||||||
/**
|
/**
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,62 @@
|
||||||
|
/*
|
||||||
|
* 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.utils;
|
||||||
|
|
||||||
|
import java.io.IOException;
|
||||||
|
import java.io.Serializable;
|
||||||
|
import java.net.InetAddress;
|
||||||
|
import java.net.Socket;
|
||||||
|
import java.rmi.server.RMIClientSocketFactory;
|
||||||
|
import java.util.Objects;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* This class is used to override the local address the JMX client calculates when trying to connect,
|
||||||
|
* which can otherwise be influenced by the system property "java.rmi.server.hostname" in strange and
|
||||||
|
* unpredictable ways.
|
||||||
|
*/
|
||||||
|
public class RMIClientSocketFactoryImpl implements RMIClientSocketFactory, Serializable
|
||||||
|
{
|
||||||
|
private final InetAddress localAddress;
|
||||||
|
|
||||||
|
public RMIClientSocketFactoryImpl(InetAddress localAddress)
|
||||||
|
{
|
||||||
|
this.localAddress = localAddress;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public Socket createSocket(String host, int port) throws IOException
|
||||||
|
{
|
||||||
|
return new Socket(localAddress, port);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean equals(Object o)
|
||||||
|
{
|
||||||
|
if (this == o) return true;
|
||||||
|
if (o == null || getClass() != o.getClass()) return false;
|
||||||
|
RMIClientSocketFactoryImpl that = (RMIClientSocketFactoryImpl) o;
|
||||||
|
return Objects.equals(localAddress, that.localAddress);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public int hashCode()
|
||||||
|
{
|
||||||
|
return Objects.hash(localAddress);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -0,0 +1,55 @@
|
||||||
|
/*
|
||||||
|
* 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.utils;
|
||||||
|
|
||||||
|
import java.lang.reflect.Field;
|
||||||
|
import java.lang.reflect.Method;
|
||||||
|
|
||||||
|
public class ReflectionUtils {
|
||||||
|
public static Field getField(Class<?> clazz, String fieldName) throws NoSuchFieldException
|
||||||
|
{
|
||||||
|
// below code works before Java 12
|
||||||
|
try
|
||||||
|
{
|
||||||
|
return clazz.getDeclaredField(fieldName);
|
||||||
|
}
|
||||||
|
catch (NoSuchFieldException e)
|
||||||
|
{
|
||||||
|
// this is mitigation for JDK 17 (https://bugs.openjdk.org/browse/JDK-8210522)
|
||||||
|
try
|
||||||
|
{
|
||||||
|
Method getDeclaredFields0 = Class.class.getDeclaredMethod("getDeclaredFields0", boolean.class);
|
||||||
|
getDeclaredFields0.setAccessible(true);
|
||||||
|
Field[] fields = (Field[]) getDeclaredFields0.invoke(clazz, false);
|
||||||
|
for (Field field : fields)
|
||||||
|
{
|
||||||
|
if (fieldName.equals(field.getName()))
|
||||||
|
{
|
||||||
|
return field;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
catch (ReflectiveOperationException ex)
|
||||||
|
{
|
||||||
|
e.addSuppressed(ex);
|
||||||
|
}
|
||||||
|
throw e;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -16,7 +16,7 @@
|
||||||
* limitations under the License.
|
* limitations under the License.
|
||||||
*/
|
*/
|
||||||
|
|
||||||
package org.apache.cassandra.distributed.shared;
|
package org.apache.cassandra.utils;
|
||||||
|
|
||||||
import java.lang.annotation.ElementType;
|
import java.lang.annotation.ElementType;
|
||||||
import java.lang.annotation.Retention;
|
import java.lang.annotation.Retention;
|
||||||
|
|
@ -70,16 +70,16 @@ import org.apache.cassandra.distributed.api.LogAction;
|
||||||
import org.apache.cassandra.distributed.api.NodeToolResult;
|
import org.apache.cassandra.distributed.api.NodeToolResult;
|
||||||
import org.apache.cassandra.distributed.api.TokenSupplier;
|
import org.apache.cassandra.distributed.api.TokenSupplier;
|
||||||
import org.apache.cassandra.distributed.shared.InstanceClassLoader;
|
import org.apache.cassandra.distributed.shared.InstanceClassLoader;
|
||||||
import org.apache.cassandra.distributed.shared.Isolated;
|
|
||||||
import org.apache.cassandra.distributed.shared.MessageFilters;
|
import org.apache.cassandra.distributed.shared.MessageFilters;
|
||||||
import org.apache.cassandra.distributed.shared.Metrics;
|
import org.apache.cassandra.distributed.shared.Metrics;
|
||||||
import org.apache.cassandra.distributed.shared.NetworkTopology;
|
import org.apache.cassandra.distributed.shared.NetworkTopology;
|
||||||
import org.apache.cassandra.distributed.shared.Shared;
|
|
||||||
import org.apache.cassandra.distributed.shared.ShutdownException;
|
import org.apache.cassandra.distributed.shared.ShutdownException;
|
||||||
import org.apache.cassandra.distributed.shared.Versions;
|
import org.apache.cassandra.distributed.shared.Versions;
|
||||||
import org.apache.cassandra.io.util.FileUtils;
|
import org.apache.cassandra.io.util.FileUtils;
|
||||||
import org.apache.cassandra.net.Verb;
|
import org.apache.cassandra.net.Verb;
|
||||||
import org.apache.cassandra.utils.FBUtilities;
|
import org.apache.cassandra.utils.FBUtilities;
|
||||||
|
import org.apache.cassandra.utils.Isolated;
|
||||||
|
import org.apache.cassandra.utils.Shared;
|
||||||
import org.apache.cassandra.utils.concurrent.SimpleCondition;
|
import org.apache.cassandra.utils.concurrent.SimpleCondition;
|
||||||
import org.reflections.Reflections;
|
import org.reflections.Reflections;
|
||||||
import org.reflections.util.ConfigurationBuilder;
|
import org.reflections.util.ConfigurationBuilder;
|
||||||
|
|
@ -165,6 +165,8 @@ public abstract class AbstractCluster<I extends IInstance> implements ICluster<I
|
||||||
private INodeProvisionStrategy.Strategy nodeProvisionStrategy = INodeProvisionStrategy.Strategy.MultipleNetworkInterfaces;
|
private INodeProvisionStrategy.Strategy nodeProvisionStrategy = INodeProvisionStrategy.Strategy.MultipleNetworkInterfaces;
|
||||||
|
|
||||||
{
|
{
|
||||||
|
// Indicate that we are running in the in-jvm dtest environment
|
||||||
|
CassandraRelevantProperties.DTEST_IS_IN_JVM_DTEST.setBoolean(true);
|
||||||
// those properties may be set for unit-test optimizations; those should not be used when running dtests
|
// those properties may be set for unit-test optimizations; those should not be used when running dtests
|
||||||
CassandraRelevantProperties.FLUSH_LOCAL_SCHEMA_CHANGES.reset();
|
CassandraRelevantProperties.FLUSH_LOCAL_SCHEMA_CHANGES.reset();
|
||||||
CassandraRelevantProperties.NON_GRACEFUL_SHUTDOWN.reset();
|
CassandraRelevantProperties.NON_GRACEFUL_SHUTDOWN.reset();
|
||||||
|
|
@ -503,6 +505,16 @@ public abstract class AbstractCluster<I extends IInstance> implements ICluster<I
|
||||||
return instanceMap.get(addr);
|
return instanceMap.get(addr);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public List<I> get(int... nodes)
|
||||||
|
{
|
||||||
|
if (nodes == null || nodes.length == 0)
|
||||||
|
throw new IllegalArgumentException("No nodes provided");
|
||||||
|
List<I> list = new ArrayList<>(nodes.length);
|
||||||
|
for (int i : nodes)
|
||||||
|
list.add(get(i));
|
||||||
|
return list;
|
||||||
|
}
|
||||||
|
|
||||||
public I getFirstRunningInstance()
|
public I getFirstRunningInstance()
|
||||||
{
|
{
|
||||||
return stream().filter(i -> !i.isShutdown()).findFirst().orElseThrow(
|
return stream().filter(i -> !i.isShutdown()).findFirst().orElseThrow(
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,87 @@
|
||||||
|
/*
|
||||||
|
* 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.impl;
|
||||||
|
|
||||||
|
import java.io.IOException;
|
||||||
|
import java.net.InetAddress;
|
||||||
|
import java.net.ServerSocket;
|
||||||
|
import java.net.SocketException;
|
||||||
|
import java.rmi.server.RMIServerSocketFactory;
|
||||||
|
import java.util.ArrayList;
|
||||||
|
import java.util.List;
|
||||||
|
import java.util.Objects;
|
||||||
|
|
||||||
|
import javax.net.ServerSocketFactory;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* This class is used to keep track of RMI servers created during a cluster creation so we can
|
||||||
|
* later close the sockets, which would otherwise be left with a thread running waiting for
|
||||||
|
* connections that would never show up as the server was otherwise closed.
|
||||||
|
*/
|
||||||
|
class CollectingRMIServerSocketFactoryImpl implements RMIServerSocketFactory
|
||||||
|
{
|
||||||
|
private final InetAddress bindAddress;
|
||||||
|
List<ServerSocket> sockets = new ArrayList<>();
|
||||||
|
|
||||||
|
public CollectingRMIServerSocketFactoryImpl(InetAddress bindAddress)
|
||||||
|
{
|
||||||
|
this.bindAddress = bindAddress;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public ServerSocket createServerSocket(int pPort) throws IOException
|
||||||
|
{
|
||||||
|
ServerSocket result = ServerSocketFactory.getDefault().createServerSocket(pPort, 0, bindAddress);
|
||||||
|
try
|
||||||
|
{
|
||||||
|
result.setReuseAddress(true);
|
||||||
|
}
|
||||||
|
catch (SocketException e)
|
||||||
|
{
|
||||||
|
result.close();
|
||||||
|
throw e;
|
||||||
|
}
|
||||||
|
sockets.add(result);
|
||||||
|
return result;
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
public void close() throws IOException
|
||||||
|
{
|
||||||
|
for (ServerSocket socket : sockets)
|
||||||
|
{
|
||||||
|
socket.close();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean equals(Object o)
|
||||||
|
{
|
||||||
|
if (this == o) return true;
|
||||||
|
if (o == null || getClass() != o.getClass()) return false;
|
||||||
|
CollectingRMIServerSocketFactoryImpl that = (CollectingRMIServerSocketFactoryImpl) o;
|
||||||
|
return Objects.equals(bindAddress, that.bindAddress);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public int hashCode()
|
||||||
|
{
|
||||||
|
return Objects.hash(bindAddress);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -52,6 +52,11 @@ public interface INodeProvisionStrategy
|
||||||
{
|
{
|
||||||
return 9041 + nodeNum;
|
return 9041 + nodeNum;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public int jmxPort(int nodeNum)
|
||||||
|
{
|
||||||
|
return 7199 + nodeNum;
|
||||||
|
}
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|
|
@ -85,6 +90,11 @@ public interface INodeProvisionStrategy
|
||||||
{
|
{
|
||||||
return 9042;
|
return 9042;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public int jmxPort(int nodeNum)
|
||||||
|
{
|
||||||
|
return 7199;
|
||||||
|
}
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
@ -96,4 +106,5 @@ public interface INodeProvisionStrategy
|
||||||
abstract String ipAddress(int nodeNum);
|
abstract String ipAddress(int nodeNum);
|
||||||
abstract int storagePort(int nodeNum);
|
abstract int storagePort(int nodeNum);
|
||||||
abstract int nativeTransportPort(int nodeNum);
|
abstract int nativeTransportPort(int nodeNum);
|
||||||
|
abstract int jmxPort(int nodeNum);
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -42,6 +42,8 @@ import java.util.concurrent.atomic.AtomicLong;
|
||||||
import javax.management.ListenerNotFoundException;
|
import javax.management.ListenerNotFoundException;
|
||||||
import javax.management.Notification;
|
import javax.management.Notification;
|
||||||
import javax.management.NotificationListener;
|
import javax.management.NotificationListener;
|
||||||
|
import javax.management.remote.JMXConnectorServer;
|
||||||
|
import javax.management.remote.rmi.RMIJRMPServerImpl;
|
||||||
|
|
||||||
import com.google.common.annotations.VisibleForTesting;
|
import com.google.common.annotations.VisibleForTesting;
|
||||||
import com.google.common.util.concurrent.Uninterruptibles;
|
import com.google.common.util.concurrent.Uninterruptibles;
|
||||||
|
|
@ -130,7 +132,10 @@ import org.apache.cassandra.utils.ByteArrayUtil;
|
||||||
import org.apache.cassandra.utils.DiagnosticSnapshotService;
|
import org.apache.cassandra.utils.DiagnosticSnapshotService;
|
||||||
import org.apache.cassandra.utils.ExecutorUtils;
|
import org.apache.cassandra.utils.ExecutorUtils;
|
||||||
import org.apache.cassandra.utils.FBUtilities;
|
import org.apache.cassandra.utils.FBUtilities;
|
||||||
|
import org.apache.cassandra.utils.JMXServerUtils;
|
||||||
import org.apache.cassandra.utils.JVMStabilityInspector;
|
import org.apache.cassandra.utils.JVMStabilityInspector;
|
||||||
|
import org.apache.cassandra.utils.MBeanWrapper;
|
||||||
|
import org.apache.cassandra.utils.RMIClientSocketFactoryImpl;
|
||||||
import org.apache.cassandra.utils.Throwables;
|
import org.apache.cassandra.utils.Throwables;
|
||||||
import org.apache.cassandra.utils.UUIDSerializer;
|
import org.apache.cassandra.utils.UUIDSerializer;
|
||||||
import org.apache.cassandra.utils.concurrent.Ref;
|
import org.apache.cassandra.utils.concurrent.Ref;
|
||||||
|
|
@ -139,6 +144,7 @@ import org.apache.cassandra.utils.progress.jmx.JMXBroadcastExecutor;
|
||||||
|
|
||||||
import static java.util.concurrent.TimeUnit.MINUTES;
|
import static java.util.concurrent.TimeUnit.MINUTES;
|
||||||
import static org.apache.cassandra.distributed.api.Feature.GOSSIP;
|
import static org.apache.cassandra.distributed.api.Feature.GOSSIP;
|
||||||
|
import static org.apache.cassandra.distributed.api.Feature.JMX;
|
||||||
import static org.apache.cassandra.distributed.api.Feature.NATIVE_PROTOCOL;
|
import static org.apache.cassandra.distributed.api.Feature.NATIVE_PROTOCOL;
|
||||||
import static org.apache.cassandra.distributed.api.Feature.NETWORK;
|
import static org.apache.cassandra.distributed.api.Feature.NETWORK;
|
||||||
import static org.apache.cassandra.distributed.impl.DistributedTestSnitch.fromCassandraInetAddressAndPort;
|
import static org.apache.cassandra.distributed.impl.DistributedTestSnitch.fromCassandraInetAddressAndPort;
|
||||||
|
|
@ -147,11 +153,19 @@ import static org.apache.cassandra.net.Verb.BATCH_STORE_REQ;
|
||||||
|
|
||||||
public class Instance extends IsolatedExecutor implements IInvokableInstance
|
public class Instance extends IsolatedExecutor implements IInvokableInstance
|
||||||
{
|
{
|
||||||
|
private static final int RMI_KEEPALIVE_TIME = 1000;
|
||||||
private Logger inInstancelogger; // Defer creation until running in the instance context
|
private Logger inInstancelogger; // Defer creation until running in the instance context
|
||||||
public final IInstanceConfig config;
|
public final IInstanceConfig config;
|
||||||
private volatile boolean initialized = false;
|
private volatile boolean initialized = false;
|
||||||
private volatile boolean internodeMessagingStarted = false;
|
private volatile boolean internodeMessagingStarted = false;
|
||||||
private final AtomicLong startedAt = new AtomicLong();
|
private final AtomicLong startedAt = new AtomicLong();
|
||||||
|
private JMXConnectorServer jmxConnectorServer;
|
||||||
|
private JMXServerUtils.JmxRegistry registry;
|
||||||
|
private RMIJRMPServerImpl jmxRmiServer;
|
||||||
|
private MBeanWrapper.InstanceMBeanWrapper wrapper;
|
||||||
|
private RMIClientSocketFactoryImpl clientSocketFactory;
|
||||||
|
private CollectingRMIServerSocketFactoryImpl serverSocketFactory;
|
||||||
|
private IsolatedJmx isolatedJmx;
|
||||||
|
|
||||||
// should never be invoked directly, so that it is instantiated on other class loader;
|
// should never be invoked directly, so that it is instantiated on other class loader;
|
||||||
// only visible for inheritance
|
// only visible for inheritance
|
||||||
|
|
@ -502,6 +516,9 @@ public class Instance extends IsolatedExecutor implements IInvokableInstance
|
||||||
config.networkTopology(), config.broadcastAddress());
|
config.networkTopology(), config.broadcastAddress());
|
||||||
DistributedTestSnitch.assign(config.networkTopology());
|
DistributedTestSnitch.assign(config.networkTopology());
|
||||||
|
|
||||||
|
if (config.has(JMX))
|
||||||
|
startJmx();
|
||||||
|
|
||||||
DatabaseDescriptor.daemonInitialization();
|
DatabaseDescriptor.daemonInitialization();
|
||||||
FileUtils.setFSErrorHandler(new DefaultFSErrorHandler());
|
FileUtils.setFSErrorHandler(new DefaultFSErrorHandler());
|
||||||
DatabaseDescriptor.createAllDirectories();
|
DatabaseDescriptor.createAllDirectories();
|
||||||
|
|
@ -624,6 +641,20 @@ public class Instance extends IsolatedExecutor implements IInvokableInstance
|
||||||
initialized = true;
|
initialized = true;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private void startJmx()
|
||||||
|
{
|
||||||
|
isolatedJmx = new IsolatedJmx(this, inInstancelogger);
|
||||||
|
isolatedJmx.startJmx();
|
||||||
|
}
|
||||||
|
|
||||||
|
private void stopJmx() throws NoSuchFieldException, InterruptedException, IllegalAccessException
|
||||||
|
{
|
||||||
|
if (config.has(JMX))
|
||||||
|
{
|
||||||
|
isolatedJmx.stopJmx();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// Update the messaging versions for all instances
|
// Update the messaging versions for all instances
|
||||||
// that have initialized their configurations.
|
// that have initialized their configurations.
|
||||||
private static void propagateMessagingVersions(ICluster cluster)
|
private static void propagateMessagingVersions(ICluster cluster)
|
||||||
|
|
@ -815,6 +846,8 @@ public class Instance extends IsolatedExecutor implements IInvokableInstance
|
||||||
|
|
||||||
// ScheduledExecutors shuts down after MessagingService, as MessagingService may issue tasks to it.
|
// ScheduledExecutors shuts down after MessagingService, as MessagingService may issue tasks to it.
|
||||||
error = parallelRun(error, executor, () -> ScheduledExecutors.shutdownAndWait(1L, MINUTES));
|
error = parallelRun(error, executor, () -> ScheduledExecutors.shutdownAndWait(1L, MINUTES));
|
||||||
|
|
||||||
|
error = parallelRun(error, executor, this::stopJmx);
|
||||||
|
|
||||||
Throwables.maybeFail(error);
|
Throwables.maybeFail(error);
|
||||||
}).apply(isolatedExecutor);
|
}).apply(isolatedExecutor);
|
||||||
|
|
|
||||||
|
|
@ -35,10 +35,10 @@ import com.vdurmont.semver4j.Semver;
|
||||||
import org.apache.cassandra.distributed.api.Feature;
|
import org.apache.cassandra.distributed.api.Feature;
|
||||||
import org.apache.cassandra.distributed.api.IInstanceConfig;
|
import org.apache.cassandra.distributed.api.IInstanceConfig;
|
||||||
import org.apache.cassandra.distributed.shared.NetworkTopology;
|
import org.apache.cassandra.distributed.shared.NetworkTopology;
|
||||||
import org.apache.cassandra.distributed.shared.Shared;
|
|
||||||
import org.apache.cassandra.distributed.upgrade.UpgradeTestBase;
|
import org.apache.cassandra.distributed.upgrade.UpgradeTestBase;
|
||||||
import org.apache.cassandra.locator.InetAddressAndPort;
|
import org.apache.cassandra.locator.InetAddressAndPort;
|
||||||
import org.apache.cassandra.locator.SimpleSeedProvider;
|
import org.apache.cassandra.locator.SimpleSeedProvider;
|
||||||
|
import org.apache.cassandra.utils.Shared;
|
||||||
|
|
||||||
@Shared
|
@Shared
|
||||||
public class InstanceConfig implements IInstanceConfig
|
public class InstanceConfig implements IInstanceConfig
|
||||||
|
|
@ -47,6 +47,8 @@ public class InstanceConfig implements IInstanceConfig
|
||||||
private static final Logger logger = LoggerFactory.getLogger(InstanceConfig.class);
|
private static final Logger logger = LoggerFactory.getLogger(InstanceConfig.class);
|
||||||
|
|
||||||
public final int num;
|
public final int num;
|
||||||
|
private final int jmxPort;
|
||||||
|
|
||||||
public int num() { return num; }
|
public int num() { return num; }
|
||||||
|
|
||||||
private final NetworkTopology networkTopology;
|
private final NetworkTopology networkTopology;
|
||||||
|
|
@ -77,7 +79,8 @@ public class InstanceConfig implements IInstanceConfig
|
||||||
String cdc_raw_directory,
|
String cdc_raw_directory,
|
||||||
String initial_token,
|
String initial_token,
|
||||||
int storage_port,
|
int storage_port,
|
||||||
int native_transport_port)
|
int native_transport_port,
|
||||||
|
int jmx_port)
|
||||||
{
|
{
|
||||||
this.num = num;
|
this.num = num;
|
||||||
this.networkTopology = networkTopology;
|
this.networkTopology = networkTopology;
|
||||||
|
|
@ -118,6 +121,7 @@ public class InstanceConfig implements IInstanceConfig
|
||||||
// legacy parameters
|
// legacy parameters
|
||||||
.forceSet("commitlog_sync_batch_window_in_ms", 1.0);
|
.forceSet("commitlog_sync_batch_window_in_ms", 1.0);
|
||||||
this.featureFlags = EnumSet.noneOf(Feature.class);
|
this.featureFlags = EnumSet.noneOf(Feature.class);
|
||||||
|
this.jmxPort = jmx_port;
|
||||||
}
|
}
|
||||||
|
|
||||||
private InstanceConfig(InstanceConfig copy)
|
private InstanceConfig(InstanceConfig copy)
|
||||||
|
|
@ -129,6 +133,7 @@ public class InstanceConfig implements IInstanceConfig
|
||||||
this.hostId = copy.hostId;
|
this.hostId = copy.hostId;
|
||||||
this.featureFlags = copy.featureFlags;
|
this.featureFlags = copy.featureFlags;
|
||||||
this.broadcastAddressAndPort = copy.broadcastAddressAndPort;
|
this.broadcastAddressAndPort = copy.broadcastAddressAndPort;
|
||||||
|
this.jmxPort = copy.jmxPort;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -174,6 +179,12 @@ public class InstanceConfig implements IInstanceConfig
|
||||||
return networkTopology().localDC(broadcastAddress());
|
return networkTopology().localDC(broadcastAddress());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public int jmxPort()
|
||||||
|
{
|
||||||
|
return this.jmxPort;
|
||||||
|
}
|
||||||
|
|
||||||
public InstanceConfig with(Feature featureFlag)
|
public InstanceConfig with(Feature featureFlag)
|
||||||
{
|
{
|
||||||
featureFlags.add(featureFlag);
|
featureFlags.add(featureFlag);
|
||||||
|
|
@ -271,7 +282,8 @@ public class InstanceConfig implements IInstanceConfig
|
||||||
String.format("%s/node%d/cdc", root, nodeNum),
|
String.format("%s/node%d/cdc", root, nodeNum),
|
||||||
token,
|
token,
|
||||||
provisionStrategy.storagePort(nodeNum),
|
provisionStrategy.storagePort(nodeNum),
|
||||||
provisionStrategy.nativeTransportPort(nodeNum));
|
provisionStrategy.nativeTransportPort(nodeNum),
|
||||||
|
provisionStrategy.jmxPort(nodeNum));
|
||||||
}
|
}
|
||||||
|
|
||||||
private static String[] datadirs(int datadirCount, File root, int nodeNum)
|
private static String[] datadirs(int datadirCount, File root, int nodeNum)
|
||||||
|
|
|
||||||
|
|
@ -53,7 +53,7 @@ public class IsolatedExecutor implements IIsolatedExecutor
|
||||||
{
|
{
|
||||||
final ExecutorService isolatedExecutor;
|
final ExecutorService isolatedExecutor;
|
||||||
private final String name;
|
private final String name;
|
||||||
private final ClassLoader classLoader;
|
final ClassLoader classLoader;
|
||||||
private final Method deserializeOnInstance;
|
private final Method deserializeOnInstance;
|
||||||
|
|
||||||
IsolatedExecutor(String name, ClassLoader classLoader)
|
IsolatedExecutor(String name, ClassLoader classLoader)
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,230 @@
|
||||||
|
/*
|
||||||
|
* 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.impl;
|
||||||
|
|
||||||
|
import java.lang.reflect.Field;
|
||||||
|
import java.net.InetAddress;
|
||||||
|
import java.net.MalformedURLException;
|
||||||
|
import java.util.HashMap;
|
||||||
|
import java.util.Iterator;
|
||||||
|
import java.util.Map;
|
||||||
|
|
||||||
|
import javax.management.remote.JMXConnector;
|
||||||
|
import javax.management.remote.JMXConnectorFactory;
|
||||||
|
import javax.management.remote.JMXConnectorServer;
|
||||||
|
import javax.management.remote.JMXServiceURL;
|
||||||
|
import javax.management.remote.rmi.RMIConnectorServer;
|
||||||
|
import javax.management.remote.rmi.RMIJRMPServerImpl;
|
||||||
|
|
||||||
|
import org.slf4j.Logger;
|
||||||
|
|
||||||
|
import org.apache.cassandra.distributed.api.IInstance;
|
||||||
|
import org.apache.cassandra.distributed.api.IInstanceConfig;
|
||||||
|
import org.apache.cassandra.utils.JMXServerUtils;
|
||||||
|
import org.apache.cassandra.utils.MBeanWrapper;
|
||||||
|
import org.apache.cassandra.utils.RMIClientSocketFactoryImpl;
|
||||||
|
import org.apache.cassandra.utils.ReflectionUtils;
|
||||||
|
import sun.rmi.transport.tcp.TCPEndpoint;
|
||||||
|
|
||||||
|
import static org.apache.cassandra.config.CassandraRelevantProperties.JAVA_RMI_DGC_LEASE_VALUE_IN_JVM_DTEST;
|
||||||
|
import static org.apache.cassandra.config.CassandraRelevantProperties.ORG_APACHE_CASSANDRA_DISABLE_MBEAN_REGISTRATION;
|
||||||
|
import static org.apache.cassandra.config.CassandraRelevantProperties.SUN_RMI_TRANSPORT_TCP_THREADKEEPALIVETIME;
|
||||||
|
import static org.apache.cassandra.distributed.api.Feature.JMX;
|
||||||
|
|
||||||
|
public class IsolatedJmx
|
||||||
|
{
|
||||||
|
private static final int RMI_KEEPALIVE_TIME = 1000;
|
||||||
|
|
||||||
|
private JMXConnectorServer jmxConnectorServer;
|
||||||
|
private JMXServerUtils.JmxRegistry registry;
|
||||||
|
private RMIJRMPServerImpl jmxRmiServer;
|
||||||
|
private MBeanWrapper.InstanceMBeanWrapper wrapper;
|
||||||
|
private RMIClientSocketFactoryImpl clientSocketFactory;
|
||||||
|
private CollectingRMIServerSocketFactoryImpl serverSocketFactory;
|
||||||
|
private Logger inInstancelogger;
|
||||||
|
private IInstanceConfig config;
|
||||||
|
|
||||||
|
public IsolatedJmx(IInstance instance, Logger inInstanceLogger) {
|
||||||
|
this.config = instance.config();
|
||||||
|
this.inInstancelogger = inInstanceLogger;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void startJmx() {
|
||||||
|
try
|
||||||
|
{
|
||||||
|
// Several RMI threads hold references to in-jvm dtest objects, and are, by default, kept
|
||||||
|
// alive for long enough (minutes) to keep classloaders from being collected.
|
||||||
|
// Set these two system properties to a low value to allow cleanup to occur fast enough
|
||||||
|
// for GC to collect our classloaders.
|
||||||
|
JAVA_RMI_DGC_LEASE_VALUE_IN_JVM_DTEST.setInt(RMI_KEEPALIVE_TIME);
|
||||||
|
SUN_RMI_TRANSPORT_TCP_THREADKEEPALIVETIME.setInt(RMI_KEEPALIVE_TIME);
|
||||||
|
ORG_APACHE_CASSANDRA_DISABLE_MBEAN_REGISTRATION.setBoolean(false);
|
||||||
|
InetAddress addr = config.broadcastAddress().getAddress();
|
||||||
|
|
||||||
|
int jmxPort = config.jmxPort();
|
||||||
|
|
||||||
|
String hostname = addr.getHostAddress();
|
||||||
|
wrapper = new MBeanWrapper.InstanceMBeanWrapper(hostname + ":" + jmxPort);
|
||||||
|
((MBeanWrapper.DelegatingMbeanWrapper) MBeanWrapper.instance).setDelegate(wrapper);
|
||||||
|
Map<String, Object> env = new HashMap<>();
|
||||||
|
|
||||||
|
serverSocketFactory = new CollectingRMIServerSocketFactoryImpl(addr);
|
||||||
|
env.put(RMIConnectorServer.RMI_SERVER_SOCKET_FACTORY_ATTRIBUTE,
|
||||||
|
serverSocketFactory);
|
||||||
|
clientSocketFactory = new RMIClientSocketFactoryImpl(addr);
|
||||||
|
env.put(RMIConnectorServer.RMI_CLIENT_SOCKET_FACTORY_ATTRIBUTE,
|
||||||
|
clientSocketFactory);
|
||||||
|
|
||||||
|
// configure the RMI registry
|
||||||
|
registry = new JMXServerUtils.JmxRegistry(jmxPort,
|
||||||
|
clientSocketFactory,
|
||||||
|
serverSocketFactory,
|
||||||
|
"jmxrmi");
|
||||||
|
|
||||||
|
// Mark the JMX server as a permanently exported object. This allows the JVM to exit with the
|
||||||
|
// server running and also exempts it from the distributed GC scheduler which otherwise would
|
||||||
|
// potentially attempt a full GC every `sun.rmi.dgc.server.gcInterval` millis (default is 3600000ms)
|
||||||
|
// For more background see:
|
||||||
|
// - CASSANDRA-2967
|
||||||
|
// - https://www.jclarity.com/2015/01/27/rmi-system-gc-unplugged/
|
||||||
|
// - https://bugs.openjdk.java.net/browse/JDK-6760712
|
||||||
|
env.put("jmx.remote.x.daemon", "true");
|
||||||
|
|
||||||
|
// Set the port used to create subsequent connections to exported objects over RMI. This simplifies
|
||||||
|
// configuration in firewalled environments, but it can't be used in conjuction with SSL sockets.
|
||||||
|
// See: CASSANDRA-7087
|
||||||
|
int rmiPort = config.jmxPort();
|
||||||
|
|
||||||
|
// We create the underlying RMIJRMPServerImpl so that we can manually bind it to the registry,
|
||||||
|
// rather then specifying a binding address in the JMXServiceURL and letting it be done automatically
|
||||||
|
// when the server is started. The reason for this is that if the registry is configured with SSL
|
||||||
|
// sockets, the JMXConnectorServer acts as its client during the binding which means it needs to
|
||||||
|
// have a truststore configured which contains the registry's certificate. Manually binding removes
|
||||||
|
// this problem.
|
||||||
|
// See CASSANDRA-12109.
|
||||||
|
jmxRmiServer = new RMIJRMPServerImpl(rmiPort, clientSocketFactory, serverSocketFactory,
|
||||||
|
env);
|
||||||
|
JMXServiceURL serviceURL = new JMXServiceURL("rmi", hostname, rmiPort);
|
||||||
|
jmxConnectorServer = new RMIConnectorServer(serviceURL, env, jmxRmiServer, wrapper.getMBeanServer());
|
||||||
|
|
||||||
|
jmxConnectorServer.start();
|
||||||
|
|
||||||
|
registry.setRemoteServerStub(jmxRmiServer.toStub());
|
||||||
|
JMXServerUtils.logJmxServiceUrl(addr, jmxPort);
|
||||||
|
waitForJmxAvailability(hostname, jmxPort, env);
|
||||||
|
}
|
||||||
|
catch (Throwable t)
|
||||||
|
{
|
||||||
|
throw new RuntimeException("Feature.JMX was enabled but could not be started.", t);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private void waitForJmxAvailability(String hostname, int rmiPort, Map<String, Object> env) throws InterruptedException, MalformedURLException
|
||||||
|
{
|
||||||
|
String url = String.format("service:jmx:rmi:///jndi/rmi://%s:%d/jmxrmi", hostname, rmiPort);
|
||||||
|
JMXServiceURL serviceURL = new JMXServiceURL(url);
|
||||||
|
int attempts = 0;
|
||||||
|
Throwable lastThrown = null;
|
||||||
|
while (attempts < 20)
|
||||||
|
{
|
||||||
|
attempts++;
|
||||||
|
try (JMXConnector ignored = JMXConnectorFactory.connect(serviceURL, env))
|
||||||
|
{
|
||||||
|
inInstancelogger.info("Connected to JMX server at {} after {} attempt(s)",
|
||||||
|
url, attempts);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
catch (MalformedURLException e)
|
||||||
|
{
|
||||||
|
throw new RuntimeException(e);
|
||||||
|
}
|
||||||
|
catch (Throwable thrown)
|
||||||
|
{
|
||||||
|
lastThrown = thrown;
|
||||||
|
}
|
||||||
|
inInstancelogger.info("Could not connect to JMX on {} after {} attempts. Will retry.", url, attempts);
|
||||||
|
Thread.sleep(1000);
|
||||||
|
}
|
||||||
|
throw new RuntimeException("Could not start JMX - unreachable after 20 attempts", lastThrown);
|
||||||
|
}
|
||||||
|
|
||||||
|
public void stopJmx() throws IllegalAccessException, NoSuchFieldException, InterruptedException
|
||||||
|
{
|
||||||
|
if (!config.has(JMX))
|
||||||
|
return;
|
||||||
|
// First, swap the mbean wrapper back to a NoOp wrapper
|
||||||
|
// This prevents later attempts to unregister mbeans from failing in Cassandra code, as we're going to
|
||||||
|
// unregister all of them here
|
||||||
|
((MBeanWrapper.DelegatingMbeanWrapper) MBeanWrapper.instance).setDelegate(new MBeanWrapper.NoOpMBeanWrapper());
|
||||||
|
try
|
||||||
|
{
|
||||||
|
wrapper.close();
|
||||||
|
}
|
||||||
|
catch (Throwable e)
|
||||||
|
{
|
||||||
|
inInstancelogger.warn("failed to close wrapper.", e);
|
||||||
|
}
|
||||||
|
try
|
||||||
|
{
|
||||||
|
jmxConnectorServer.stop();
|
||||||
|
}
|
||||||
|
catch (Throwable e)
|
||||||
|
{
|
||||||
|
inInstancelogger.warn("failed to close jmxConnectorServer.", e);
|
||||||
|
}
|
||||||
|
try
|
||||||
|
{
|
||||||
|
registry.close();
|
||||||
|
}
|
||||||
|
catch (Throwable e)
|
||||||
|
{
|
||||||
|
inInstancelogger.warn("failed to close registry.", e);
|
||||||
|
}
|
||||||
|
try
|
||||||
|
{
|
||||||
|
serverSocketFactory.close();
|
||||||
|
}
|
||||||
|
catch (Throwable e)
|
||||||
|
{
|
||||||
|
inInstancelogger.warn("failed to close serverSocketFactory.", e);
|
||||||
|
}
|
||||||
|
// The TCPEndpoint class holds references to a class in the in-jvm dtest framework
|
||||||
|
// which transitively has a reference to the InstanceClassLoader, so we need to
|
||||||
|
// make sure to remove the reference to them when the instance is shutting down
|
||||||
|
clearMapField(TCPEndpoint.class, null, "localEndpoints");
|
||||||
|
Thread.sleep(2 * RMI_KEEPALIVE_TIME); // Double the keep-alive time to give Distributed GC some time to clean up
|
||||||
|
}
|
||||||
|
|
||||||
|
private <K, V> void clearMapField(Class<?> clazz, Object instance, String mapName)
|
||||||
|
throws IllegalAccessException, NoSuchFieldException {
|
||||||
|
Field mapField = ReflectionUtils.getField(clazz, mapName);
|
||||||
|
mapField.setAccessible(true);
|
||||||
|
Map<K, V> map = (Map<K, V>) mapField.get(instance);
|
||||||
|
// Because multiple instances can be shutting down at once,
|
||||||
|
// synchronize on the map to avoid ConcurrentModificationException
|
||||||
|
synchronized (map)
|
||||||
|
{
|
||||||
|
for (Iterator<Map.Entry<K, V>> it = map.entrySet().iterator(); it.hasNext(); )
|
||||||
|
{
|
||||||
|
it.next();
|
||||||
|
it.remove();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -52,6 +52,7 @@ import org.apache.cassandra.distributed.impl.AbstractCluster;
|
||||||
import org.apache.cassandra.distributed.impl.InstanceConfig;
|
import org.apache.cassandra.distributed.impl.InstanceConfig;
|
||||||
import org.apache.cassandra.service.StorageService;
|
import org.apache.cassandra.service.StorageService;
|
||||||
import org.apache.cassandra.utils.FBUtilities;
|
import org.apache.cassandra.utils.FBUtilities;
|
||||||
|
import org.apache.cassandra.utils.Isolated;
|
||||||
|
|
||||||
import static com.google.common.util.concurrent.Uninterruptibles.sleepUninterruptibly;
|
import static com.google.common.util.concurrent.Uninterruptibles.sleepUninterruptibly;
|
||||||
import static org.apache.cassandra.config.CassandraRelevantProperties.BOOTSTRAP_SCHEMA_DELAY_MS;
|
import static org.apache.cassandra.config.CassandraRelevantProperties.BOOTSTRAP_SCHEMA_DELAY_MS;
|
||||||
|
|
|
||||||
|
|
@ -36,7 +36,7 @@ import org.apache.cassandra.distributed.api.SimpleQueryResult;
|
||||||
import org.apache.cassandra.distributed.api.TokenSupplier;
|
import org.apache.cassandra.distributed.api.TokenSupplier;
|
||||||
import org.apache.cassandra.distributed.shared.Byteman;
|
import org.apache.cassandra.distributed.shared.Byteman;
|
||||||
import org.apache.cassandra.distributed.shared.NetworkTopology;
|
import org.apache.cassandra.distributed.shared.NetworkTopology;
|
||||||
import org.apache.cassandra.distributed.shared.Shared;
|
import org.apache.cassandra.utils.Shared;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Replaces python dtest bootstrap_test.py::TestBootstrap::test_bootstrap_binary_disabled
|
* Replaces python dtest bootstrap_test.py::TestBootstrap::test_bootstrap_binary_disabled
|
||||||
|
|
|
||||||
|
|
@ -28,19 +28,27 @@ import java.text.SimpleDateFormat;
|
||||||
import java.time.Instant;
|
import java.time.Instant;
|
||||||
import java.util.function.Consumer;
|
import java.util.function.Consumer;
|
||||||
import javax.management.MBeanServer;
|
import javax.management.MBeanServer;
|
||||||
|
import javax.management.MBeanServerConnection;
|
||||||
|
import javax.management.remote.JMXConnector;
|
||||||
|
|
||||||
|
import org.junit.Assert;
|
||||||
import org.junit.Ignore;
|
import org.junit.Ignore;
|
||||||
import org.junit.Test;
|
import org.junit.Test;
|
||||||
|
|
||||||
import com.sun.management.HotSpotDiagnosticMXBean;
|
import com.sun.management.HotSpotDiagnosticMXBean;
|
||||||
import org.apache.cassandra.distributed.Cluster;
|
import org.apache.cassandra.distributed.Cluster;
|
||||||
import org.apache.cassandra.distributed.api.ConsistencyLevel;
|
import org.apache.cassandra.distributed.api.ConsistencyLevel;
|
||||||
|
import org.apache.cassandra.distributed.api.Feature;
|
||||||
import org.apache.cassandra.distributed.api.IInstanceConfig;
|
import org.apache.cassandra.distributed.api.IInstanceConfig;
|
||||||
|
import org.apache.cassandra.distributed.api.IInvokableInstance;
|
||||||
|
import org.apache.cassandra.distributed.shared.JMXUtil;
|
||||||
import org.apache.cassandra.utils.SigarLibrary;
|
import org.apache.cassandra.utils.SigarLibrary;
|
||||||
|
|
||||||
import static org.apache.cassandra.distributed.api.Feature.GOSSIP;
|
import static org.apache.cassandra.distributed.api.Feature.GOSSIP;
|
||||||
|
import static org.apache.cassandra.distributed.api.Feature.JMX;
|
||||||
import static org.apache.cassandra.distributed.api.Feature.NATIVE_PROTOCOL;
|
import static org.apache.cassandra.distributed.api.Feature.NATIVE_PROTOCOL;
|
||||||
import static org.apache.cassandra.distributed.api.Feature.NETWORK;
|
import static org.apache.cassandra.distributed.api.Feature.NETWORK;
|
||||||
|
import static org.hamcrest.Matchers.startsWith;
|
||||||
|
|
||||||
/* Resource Leak Test - useful when tracking down issues with in-JVM framework cleanup.
|
/* Resource Leak Test - useful when tracking down issues with in-JVM framework cleanup.
|
||||||
* All objects referencing the InstanceClassLoader need to be garbage collected or
|
* All objects referencing the InstanceClassLoader need to be garbage collected or
|
||||||
|
|
@ -138,6 +146,11 @@ public class ResourceLeakTest extends TestBaseImpl
|
||||||
}
|
}
|
||||||
|
|
||||||
void doTest(int numClusterNodes, Consumer<IInstanceConfig> updater) throws Throwable
|
void doTest(int numClusterNodes, Consumer<IInstanceConfig> updater) throws Throwable
|
||||||
|
{
|
||||||
|
doTest(numClusterNodes, updater, ignored -> {});
|
||||||
|
}
|
||||||
|
|
||||||
|
void doTest(int numClusterNodes, Consumer<IInstanceConfig> updater, Consumer<Cluster> actionToPerform) throws Throwable
|
||||||
{
|
{
|
||||||
for (int loop = 0; loop < numTestLoops; loop++)
|
for (int loop = 0; loop < numTestLoops; loop++)
|
||||||
{
|
{
|
||||||
|
|
@ -149,6 +162,7 @@ public class ResourceLeakTest extends TestBaseImpl
|
||||||
cluster.schemaChange("CREATE TABLE " + KEYSPACE + "." + tableName + " (pk int, ck int, v int, PRIMARY KEY (pk, ck))");
|
cluster.schemaChange("CREATE TABLE " + KEYSPACE + "." + tableName + " (pk int, ck int, v int, PRIMARY KEY (pk, ck))");
|
||||||
cluster.coordinator(1).execute("INSERT INTO " + KEYSPACE + "." + tableName + "(pk,ck,v) VALUES (0,0,0)", ConsistencyLevel.ALL);
|
cluster.coordinator(1).execute("INSERT INTO " + KEYSPACE + "." + tableName + "(pk,ck,v) VALUES (0,0,0)", ConsistencyLevel.ALL);
|
||||||
cluster.get(1).flush(KEYSPACE);
|
cluster.get(1).flush(KEYSPACE);
|
||||||
|
actionToPerform.accept(cluster);
|
||||||
if (dumpEveryLoop)
|
if (dumpEveryLoop)
|
||||||
{
|
{
|
||||||
dumpResources(String.format("loop%03d", loop));
|
dumpResources(String.format("loop%03d", loop));
|
||||||
|
|
@ -207,4 +221,50 @@ public class ResourceLeakTest extends TestBaseImpl
|
||||||
}
|
}
|
||||||
dumpResources("final-native");
|
dumpResources("final-native");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void looperJmxTest() throws Throwable
|
||||||
|
{
|
||||||
|
doTest(1, config -> config.with(JMX), cluster -> {
|
||||||
|
// NOTE: At some point, the hostname of the broadcastAddress can be resolved
|
||||||
|
// and then the `getHostString`, which would otherwise return the IP address,
|
||||||
|
// starts returning `localhost` - use `.getAddress().getHostAddress()` to work around this.
|
||||||
|
for (IInvokableInstance instance:cluster.get(1, cluster.size()))
|
||||||
|
{
|
||||||
|
IInstanceConfig config = instance.config();
|
||||||
|
try (JMXConnector jmxc = JMXUtil.getJmxConnector(config))
|
||||||
|
{
|
||||||
|
MBeanServerConnection mbsc = jmxc.getMBeanServerConnection();
|
||||||
|
// instances get their default domain set to their IP address, so us it
|
||||||
|
// to check that we are actually connecting to the correct instance
|
||||||
|
String defaultDomain = mbsc.getDefaultDomain();
|
||||||
|
Assert.assertThat(defaultDomain, startsWith(JMXUtil.getJmxHost(config) + ":" + config.jmxPort()));
|
||||||
|
}
|
||||||
|
catch (IOException e)
|
||||||
|
{
|
||||||
|
throw new RuntimeException(e);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
});
|
||||||
|
if (forceCollection)
|
||||||
|
{
|
||||||
|
System.runFinalization();
|
||||||
|
System.gc();
|
||||||
|
Thread.sleep(finalWaitMillis);
|
||||||
|
}
|
||||||
|
dumpResources("final-jmx");
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void looperEverythingTest() throws Throwable
|
||||||
|
{
|
||||||
|
doTest(1, config -> config.with(Feature.values()));
|
||||||
|
if (forceCollection)
|
||||||
|
{
|
||||||
|
System.runFinalization();
|
||||||
|
System.gc();
|
||||||
|
Thread.sleep(finalWaitMillis);
|
||||||
|
}
|
||||||
|
dumpResources("final-everything");
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,113 @@
|
||||||
|
/*
|
||||||
|
* 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.jmx;
|
||||||
|
|
||||||
|
import java.io.IOException;
|
||||||
|
import java.util.HashSet;
|
||||||
|
import java.util.Set;
|
||||||
|
import javax.management.MBeanServerConnection;
|
||||||
|
import javax.management.remote.JMXConnector;
|
||||||
|
|
||||||
|
import org.junit.Assert;
|
||||||
|
import org.junit.Test;
|
||||||
|
|
||||||
|
import org.apache.cassandra.distributed.Cluster;
|
||||||
|
import org.apache.cassandra.distributed.api.Feature;
|
||||||
|
import org.apache.cassandra.distributed.api.IInstanceConfig;
|
||||||
|
import org.apache.cassandra.distributed.api.IInvokableInstance;
|
||||||
|
import org.apache.cassandra.distributed.impl.INodeProvisionStrategy;
|
||||||
|
import org.apache.cassandra.distributed.shared.JMXUtil;
|
||||||
|
import org.apache.cassandra.distributed.test.TestBaseImpl;
|
||||||
|
|
||||||
|
import static org.hamcrest.Matchers.startsWith;
|
||||||
|
|
||||||
|
public class JMXFeatureTest extends TestBaseImpl
|
||||||
|
{
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Test the in-jvm dtest JMX feature.
|
||||||
|
* - Create a cluster with multiple JMX servers, one per instance
|
||||||
|
* - Test that when connecting, we get the correct MBeanServer by checking the default domain, which is set to the IP of the instance
|
||||||
|
* - Run the test multiple times to ensure cleanup of the JMX servers is complete so the next test can run successfully using the same host/port.
|
||||||
|
*
|
||||||
|
* @throws Exception
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
public void testMultipleNetworkInterfacesProvisioning() throws Exception
|
||||||
|
{
|
||||||
|
int iterations = 2; // Make sure the JMX infrastructure all cleans up properly by running this multiple times.
|
||||||
|
Set<String> allInstances = new HashSet<>();
|
||||||
|
for (int i = 0; i < iterations; i++)
|
||||||
|
{
|
||||||
|
try (Cluster cluster = Cluster.build(2)
|
||||||
|
.withNodeProvisionStrategy(INodeProvisionStrategy.Strategy.MultipleNetworkInterfaces)
|
||||||
|
.withConfig(c -> c.with(Feature.values())).start())
|
||||||
|
{
|
||||||
|
Set<String> instancesContacted = new HashSet<>();
|
||||||
|
for (IInvokableInstance instance : cluster.get(1, 2))
|
||||||
|
{
|
||||||
|
testInstance(instancesContacted, instance);
|
||||||
|
}
|
||||||
|
Assert.assertEquals("Should have connected with both JMX instances.", 2, instancesContacted.size());
|
||||||
|
allInstances.addAll(instancesContacted);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Assert.assertEquals("Each instance from each cluster should have been unique", iterations * 2, allInstances.size());
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void testOneNetworkInterfaceProvisioning() throws Exception
|
||||||
|
{
|
||||||
|
int iterations = 2; // Make sure the JMX infrastructure all cleans up properly by running this multiple times.
|
||||||
|
Set<String> allInstances = new HashSet<>();
|
||||||
|
for (int i = 0; i < iterations; i++)
|
||||||
|
{
|
||||||
|
try (Cluster cluster = Cluster.build(2)
|
||||||
|
.withNodeProvisionStrategy(INodeProvisionStrategy.Strategy.OneNetworkInterface)
|
||||||
|
.withConfig(c -> c.with(Feature.values())).start())
|
||||||
|
{
|
||||||
|
Set<String> instancesContacted = new HashSet<>();
|
||||||
|
for (IInvokableInstance instance : cluster.get(1, 2))
|
||||||
|
{
|
||||||
|
testInstance(instancesContacted, instance);
|
||||||
|
}
|
||||||
|
Assert.assertEquals("Should have connected with both JMX instances.", 2, instancesContacted.size());
|
||||||
|
allInstances.addAll(instancesContacted);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Assert.assertEquals("Each instance from each cluster should have been unique", iterations * 2, allInstances.size());
|
||||||
|
}
|
||||||
|
|
||||||
|
private void testInstance(Set<String> instancesContacted, IInvokableInstance instance) throws IOException
|
||||||
|
{
|
||||||
|
// NOTE: At some point, the hostname of the broadcastAddress can be resolved
|
||||||
|
// and then the `getHostString`, which would otherwise return the IP address,
|
||||||
|
// starts returning `localhost` - use `.getAddress().getHostAddress()` to work around this.
|
||||||
|
IInstanceConfig config = instance.config();
|
||||||
|
try (JMXConnector jmxc = JMXUtil.getJmxConnector(config))
|
||||||
|
{
|
||||||
|
MBeanServerConnection mbsc = jmxc.getMBeanServerConnection();
|
||||||
|
// instances get their default domain set to their IP address, so us it
|
||||||
|
// to check that we are actually connecting to the correct instance
|
||||||
|
String defaultDomain = mbsc.getDefaultDomain();
|
||||||
|
instancesContacted.add(defaultDomain);
|
||||||
|
Assert.assertThat(defaultDomain, startsWith(JMXUtil.getJmxHost(config) + ":" + config.jmxPort()));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -0,0 +1,136 @@
|
||||||
|
/*
|
||||||
|
* 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.jmx;
|
||||||
|
|
||||||
|
import java.util.ArrayList;
|
||||||
|
import java.util.Collections;
|
||||||
|
import java.util.List;
|
||||||
|
import java.util.Set;
|
||||||
|
import java.util.TreeSet;
|
||||||
|
import javax.management.JMRuntimeException;
|
||||||
|
import javax.management.MBeanAttributeInfo;
|
||||||
|
import javax.management.MBeanInfo;
|
||||||
|
import javax.management.MBeanOperationInfo;
|
||||||
|
import javax.management.MBeanServerConnection;
|
||||||
|
import javax.management.ObjectName;
|
||||||
|
import javax.management.remote.JMXConnector;
|
||||||
|
import javax.management.remote.JMXConnectorFactory;
|
||||||
|
import javax.management.remote.JMXServiceURL;
|
||||||
|
|
||||||
|
import com.google.common.collect.ImmutableSet;
|
||||||
|
import org.junit.Test;
|
||||||
|
|
||||||
|
import org.apache.cassandra.distributed.Cluster;
|
||||||
|
import org.apache.cassandra.distributed.api.Feature;
|
||||||
|
import org.apache.cassandra.distributed.api.IInvokableInstance;
|
||||||
|
import org.apache.cassandra.distributed.test.TestBaseImpl;
|
||||||
|
|
||||||
|
public class JMXGetterCheckTest extends TestBaseImpl
|
||||||
|
{
|
||||||
|
private static final Set<String> IGNORE_ATTRIBUTES = ImmutableSet.of(
|
||||||
|
"org.apache.cassandra.net:type=MessagingService:BackPressurePerHost" // throws unsupported saying the feature was removed... dropped in CASSANDRA-15375
|
||||||
|
);
|
||||||
|
private static final Set<String> IGNORE_OPERATIONS = ImmutableSet.of(
|
||||||
|
"org.apache.cassandra.db:type=StorageService:stopDaemon", // halts the instance, which then causes the JVM to exit
|
||||||
|
"org.apache.cassandra.db:type=StorageService:drain", // don't drain, it stops things which can cause other APIs to be unstable as we are in a stopped state
|
||||||
|
"org.apache.cassandra.db:type=StorageService:stopGossiping", // if we stop gossip this can cause other issues, so avoid
|
||||||
|
"org.apache.cassandra.db:type=StorageService:resetLocalSchema" // this will fail when there are no other nodes which can serve schema
|
||||||
|
);
|
||||||
|
|
||||||
|
public static final String JMX_SERVICE_URL_FMT = "service:jmx:rmi:///jndi/rmi://%s:%d/jmxrmi";
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void testGetters() throws Exception
|
||||||
|
{
|
||||||
|
try (Cluster cluster = Cluster.build(1).withConfig(c -> c.with(Feature.values())).start())
|
||||||
|
{
|
||||||
|
IInvokableInstance instance = cluster.get(1);
|
||||||
|
|
||||||
|
String jmxHost = instance.config().broadcastAddress().getAddress().getHostAddress();
|
||||||
|
String url = String.format(JMX_SERVICE_URL_FMT, jmxHost, instance.config().jmxPort());
|
||||||
|
List<Named> errors = new ArrayList<>();
|
||||||
|
try (JMXConnector jmxc = JMXConnectorFactory.connect(new JMXServiceURL(url), null))
|
||||||
|
{
|
||||||
|
MBeanServerConnection mbsc = jmxc.getMBeanServerConnection();
|
||||||
|
Set<ObjectName> metricNames = new TreeSet<>(mbsc.queryNames(null, null));
|
||||||
|
for (ObjectName name : metricNames)
|
||||||
|
{
|
||||||
|
if (!name.getDomain().startsWith("org.apache.cassandra"))
|
||||||
|
continue;
|
||||||
|
MBeanInfo info = mbsc.getMBeanInfo(name);
|
||||||
|
for (MBeanAttributeInfo a : info.getAttributes())
|
||||||
|
{
|
||||||
|
String fqn = String.format("%s:%s", name, a.getName());
|
||||||
|
if (!a.isReadable() || IGNORE_ATTRIBUTES.contains(fqn))
|
||||||
|
continue;
|
||||||
|
try
|
||||||
|
{
|
||||||
|
mbsc.getAttribute(name, a.getName());
|
||||||
|
}
|
||||||
|
catch (JMRuntimeException e)
|
||||||
|
{
|
||||||
|
errors.add(new Named(String.format("Attribute %s", fqn), e.getCause()));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
for (MBeanOperationInfo o : info.getOperations())
|
||||||
|
{
|
||||||
|
String fqn = String.format("%s:%s", name, o.getName());
|
||||||
|
if (o.getSignature().length != 0 || IGNORE_OPERATIONS.contains(fqn))
|
||||||
|
continue;
|
||||||
|
try
|
||||||
|
{
|
||||||
|
mbsc.invoke(name, o.getName(), new Object[0], new String[0]);
|
||||||
|
}
|
||||||
|
catch (JMRuntimeException e)
|
||||||
|
{
|
||||||
|
errors.add(new Named(String.format("Operation %s", fqn), e.getCause()));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if (!errors.isEmpty())
|
||||||
|
{
|
||||||
|
AssertionError root = new AssertionError();
|
||||||
|
errors.forEach(root::addSuppressed);
|
||||||
|
throw root;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* This class is meant to make new errors easier to read, by adding the JMX endpoint, and cleaning up the unneeded JMX/Reflection logic cluttering the stacktrace
|
||||||
|
*/
|
||||||
|
private static class Named extends RuntimeException
|
||||||
|
{
|
||||||
|
public Named(String msg, Throwable cause)
|
||||||
|
{
|
||||||
|
super(msg + "\nCaused by: " + cause.getClass().getCanonicalName() + ": " + cause.getMessage(), cause.getCause());
|
||||||
|
StackTraceElement[] stack = cause.getStackTrace();
|
||||||
|
List<StackTraceElement> copy = new ArrayList<>();
|
||||||
|
for (StackTraceElement s : stack)
|
||||||
|
{
|
||||||
|
if (!s.getClassName().startsWith("org.apache.cassandra"))
|
||||||
|
break;
|
||||||
|
copy.add(s);
|
||||||
|
}
|
||||||
|
Collections.reverse(copy);
|
||||||
|
setStackTrace(copy.toArray(new StackTraceElement[0]));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -29,7 +29,9 @@ import java.util.concurrent.ConcurrentMap;
|
||||||
import java.util.stream.Collectors;
|
import java.util.stream.Collectors;
|
||||||
import javax.management.InstanceAlreadyExistsException;
|
import javax.management.InstanceAlreadyExistsException;
|
||||||
import javax.management.InstanceNotFoundException;
|
import javax.management.InstanceNotFoundException;
|
||||||
|
import javax.management.MBeanServer;
|
||||||
import javax.management.ObjectName;
|
import javax.management.ObjectName;
|
||||||
|
import javax.management.QueryExp;
|
||||||
|
|
||||||
import com.google.common.collect.ImmutableSet;
|
import com.google.common.collect.ImmutableSet;
|
||||||
import org.junit.Assert;
|
import org.junit.Assert;
|
||||||
|
|
@ -56,7 +58,6 @@ public class TableMetricTest extends TestBaseImpl
|
||||||
MBEAN_REGISTRATION_CLASS.setString(MapMBeanWrapper.class.getName());
|
MBEAN_REGISTRATION_CLASS.setString(MapMBeanWrapper.class.getName());
|
||||||
IS_DISABLED_MBEAN_REGISTRATION.setBoolean(false);
|
IS_DISABLED_MBEAN_REGISTRATION.setBoolean(false);
|
||||||
}
|
}
|
||||||
|
|
||||||
private static volatile Map<String, Collection<String>> SYSTEM_TABLES = null;
|
private static volatile Map<String, Collection<String>> SYSTEM_TABLES = null;
|
||||||
private static Set<String> TABLE_METRIC_NAMES = ImmutableSet.of("WriteLatency");
|
private static Set<String> TABLE_METRIC_NAMES = ImmutableSet.of("WriteLatency");
|
||||||
|
|
||||||
|
|
@ -164,7 +165,7 @@ public class TableMetricTest extends TestBaseImpl
|
||||||
{
|
{
|
||||||
inst.runOnInstance(() -> {
|
inst.runOnInstance(() -> {
|
||||||
// cast only to make sure it linked properly
|
// cast only to make sure it linked properly
|
||||||
MapMBeanWrapper mbeans = (MapMBeanWrapper) MBeanWrapper.instance;
|
MapMBeanWrapper mbeans = getMapMBeanWrapper();
|
||||||
Assert.assertTrue("Unable to find table mbean for " + keyspace + "." + table,
|
Assert.assertTrue("Unable to find table mbean for " + keyspace + "." + table,
|
||||||
mbeans.isRegistered(ColumnFamilyStore.getTableMBeanName(keyspace, table, false)));
|
mbeans.isRegistered(ColumnFamilyStore.getTableMBeanName(keyspace, table, false)));
|
||||||
Assert.assertTrue("Unable to find column family mbean for " + keyspace + "." + table,
|
Assert.assertTrue("Unable to find column family mbean for " + keyspace + "." + table,
|
||||||
|
|
@ -176,7 +177,7 @@ public class TableMetricTest extends TestBaseImpl
|
||||||
{
|
{
|
||||||
inst.runOnInstance(() -> {
|
inst.runOnInstance(() -> {
|
||||||
// cast only to make sure it linked properly
|
// cast only to make sure it linked properly
|
||||||
MapMBeanWrapper mbeans = (MapMBeanWrapper) MBeanWrapper.instance;
|
MapMBeanWrapper mbeans = getMapMBeanWrapper();
|
||||||
Assert.assertFalse("Found table mbean for " + keyspace + "." + table,
|
Assert.assertFalse("Found table mbean for " + keyspace + "." + table,
|
||||||
mbeans.isRegistered(ColumnFamilyStore.getTableMBeanName(keyspace, table, false)));
|
mbeans.isRegistered(ColumnFamilyStore.getTableMBeanName(keyspace, table, false)));
|
||||||
Assert.assertFalse("Found column family mbean for " + keyspace + "." + table,
|
Assert.assertFalse("Found column family mbean for " + keyspace + "." + table,
|
||||||
|
|
@ -188,7 +189,7 @@ public class TableMetricTest extends TestBaseImpl
|
||||||
{
|
{
|
||||||
inst.runOnInstance(() -> {
|
inst.runOnInstance(() -> {
|
||||||
// cast only to make sure it linked properly
|
// cast only to make sure it linked properly
|
||||||
MapMBeanWrapper mbeans = (MapMBeanWrapper) MBeanWrapper.instance;
|
MapMBeanWrapper mbeans = getMapMBeanWrapper();
|
||||||
String mbean = getTableMetricName(keyspace, table, name);
|
String mbean = getTableMetricName(keyspace, table, name);
|
||||||
Assert.assertTrue("Unable to find metric " + name + " for " + keyspace + "." + table, mbeans.isRegistered(mbean));
|
Assert.assertTrue("Unable to find metric " + name + " for " + keyspace + "." + table, mbeans.isRegistered(mbean));
|
||||||
|
|
||||||
|
|
@ -202,7 +203,7 @@ public class TableMetricTest extends TestBaseImpl
|
||||||
{
|
{
|
||||||
inst.runOnInstance(() -> {
|
inst.runOnInstance(() -> {
|
||||||
// cast only to make sure it linked properly
|
// cast only to make sure it linked properly
|
||||||
MapMBeanWrapper mbeans = (MapMBeanWrapper) MBeanWrapper.instance;
|
MapMBeanWrapper mbeans = getMapMBeanWrapper();
|
||||||
String mbean = getTableMetricName(keyspace, table, name);
|
String mbean = getTableMetricName(keyspace, table, name);
|
||||||
Assert.assertFalse("Found metric " + name + " for " + keyspace + "." + table, mbeans.isRegistered(mbean));
|
Assert.assertFalse("Found metric " + name + " for " + keyspace + "." + table, mbeans.isRegistered(mbean));
|
||||||
|
|
||||||
|
|
@ -225,7 +226,7 @@ public class TableMetricTest extends TestBaseImpl
|
||||||
{
|
{
|
||||||
inst.runOnInstance(() -> {
|
inst.runOnInstance(() -> {
|
||||||
// cast only to make sure it linked properly
|
// cast only to make sure it linked properly
|
||||||
MapMBeanWrapper mbeans = (MapMBeanWrapper) MBeanWrapper.instance;
|
MapMBeanWrapper mbeans = getMapMBeanWrapper();
|
||||||
|
|
||||||
String keyspaceMBean = getKeyspaceMetricName(keyspace, name);
|
String keyspaceMBean = getKeyspaceMetricName(keyspace, name);
|
||||||
Assert.assertFalse("Found keyspace metric " + keyspaceMBean + " for " + keyspace, mbeans.isRegistered(keyspaceMBean));
|
Assert.assertFalse("Found keyspace metric " + keyspaceMBean + " for " + keyspace, mbeans.isRegistered(keyspaceMBean));
|
||||||
|
|
@ -242,6 +243,10 @@ public class TableMetricTest extends TestBaseImpl
|
||||||
return String.format("org.apache.cassandra.metrics:type=Table,keyspace=%s,scope=%s,name=%s", keyspace, table, name);
|
return String.format("org.apache.cassandra.metrics:type=Table,keyspace=%s,scope=%s,name=%s", keyspace, table, name);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private static MapMBeanWrapper getMapMBeanWrapper()
|
||||||
|
{
|
||||||
|
return (MapMBeanWrapper) ((MBeanWrapper.DelegatingMbeanWrapper)MBeanWrapper.instance).getDelegate();
|
||||||
|
}
|
||||||
public static final class MapMBeanWrapper implements MBeanWrapper
|
public static final class MapMBeanWrapper implements MBeanWrapper
|
||||||
{
|
{
|
||||||
private final ConcurrentMap<ObjectName, Object> map = new ConcurrentHashMap<>();
|
private final ConcurrentMap<ObjectName, Object> map = new ConcurrentHashMap<>();
|
||||||
|
|
@ -267,5 +272,17 @@ public class TableMetricTest extends TestBaseImpl
|
||||||
if (previous == null)
|
if (previous == null)
|
||||||
onException.handler.accept(new InstanceNotFoundException("MBean " + mbeanName + " was not found"));
|
onException.handler.accept(new InstanceNotFoundException("MBean " + mbeanName + " was not found"));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public Set<ObjectName> queryNames(ObjectName name, QueryExp query)
|
||||||
|
{
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public MBeanServer getMBeanServer()
|
||||||
|
{
|
||||||
|
return null;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -30,7 +30,7 @@ import org.apache.cassandra.distributed.UpgradeableCluster;
|
||||||
import org.apache.cassandra.distributed.api.ConsistencyLevel;
|
import org.apache.cassandra.distributed.api.ConsistencyLevel;
|
||||||
import org.apache.cassandra.distributed.api.Feature;
|
import org.apache.cassandra.distributed.api.Feature;
|
||||||
import org.apache.cassandra.distributed.api.IUpgradeableInstance;
|
import org.apache.cassandra.distributed.api.IUpgradeableInstance;
|
||||||
import org.apache.cassandra.distributed.shared.Shared;
|
import org.apache.cassandra.utils.Shared;
|
||||||
import org.awaitility.Awaitility;
|
import org.awaitility.Awaitility;
|
||||||
|
|
||||||
import static org.apache.cassandra.distributed.shared.AssertUtils.*;
|
import static org.apache.cassandra.distributed.shared.AssertUtils.*;
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue