From d348f353220b3b094da3093fcc5ca50f3e8b6f75 Mon Sep 17 00:00:00 2001 From: David Webb Date: Wed, 22 Jan 2014 13:44:28 -0500 Subject: [PATCH] DATACASS-38 - Cluster Connection Listener Added Host.StateListener Added LatencyTracker Added SSL/SSLOptions --- .../config/CassandraClusterFactoryBean.java | 52 ++++++++++ .../config/xml/CassandraClusterParser.java | 30 +++++- .../cassandra/config/spring-cassandra-1.0.xsd | 99 +++++++++++++++---- ...tractEmbeddedCassandraIntegrationTest.java | 2 +- .../config/xml/TestHostStateListener.java | 52 ++++++++++ .../config/xml/TestLatencyTracker.java | 37 +++++++ ...NamespaceCreatingXmlConfigTest-context.xml | 8 +- .../integration/config/xml/ppncxct.properties | 1 + 8 files changed, 255 insertions(+), 26 deletions(-) create mode 100644 spring-cassandra/src/test/java/org/springframework/cassandra/test/integration/config/xml/TestHostStateListener.java create mode 100644 spring-cassandra/src/test/java/org/springframework/cassandra/test/integration/config/xml/TestLatencyTracker.java diff --git a/spring-cassandra/src/main/java/org/springframework/cassandra/config/CassandraClusterFactoryBean.java b/spring-cassandra/src/main/java/org/springframework/cassandra/config/CassandraClusterFactoryBean.java index 69453baec..acb7375af 100644 --- a/spring-cassandra/src/main/java/org/springframework/cassandra/config/CassandraClusterFactoryBean.java +++ b/spring-cassandra/src/main/java/org/springframework/cassandra/config/CassandraClusterFactoryBean.java @@ -39,8 +39,11 @@ import org.springframework.util.StringUtils; import com.datastax.driver.core.AuthProvider; import com.datastax.driver.core.Cluster; +import com.datastax.driver.core.Host; +import com.datastax.driver.core.LatencyTracker; import com.datastax.driver.core.PoolingOptions; import com.datastax.driver.core.ProtocolOptions.Compression; +import com.datastax.driver.core.SSLOptions; import com.datastax.driver.core.Session; import com.datastax.driver.core.SocketOptions; import com.datastax.driver.core.policies.LoadBalancingPolicy; @@ -61,6 +64,7 @@ public class CassandraClusterFactoryBean implements FactoryBean, Initia public static final boolean DEFAULT_METRICS_ENABLED = true; public static final boolean DEFAULT_DEFERRED_INITIALIZATION = false; public static final boolean DEFAULT_JMX_REPORTING_ENABLED = true; + public static final boolean DEFAULT_SSL_ENABLED = false; public static final int DEFAULT_PORT = 9042; protected static final Logger log = LoggerFactory.getLogger(CassandraClusterFactoryBean.class); @@ -84,6 +88,10 @@ public class CassandraClusterFactoryBean implements FactoryBean, Initia private boolean metricsEnabled = DEFAULT_METRICS_ENABLED; private boolean deferredInitialization = DEFAULT_DEFERRED_INITIALIZATION; private boolean jmxReportingEnabled = DEFAULT_JMX_REPORTING_ENABLED; + private boolean sslEnabled = DEFAULT_SSL_ENABLED; + private SSLOptions sslOptions; + private Host.StateListener hostStateListener; + private LatencyTracker latencyTracker; private Set> keyspaceSpecifications = new HashSet>(); private List keyspaceCreations = new ArrayList(); private List keyspaceDrops = new ArrayList(); @@ -167,8 +175,24 @@ public class CassandraClusterFactoryBean implements FactoryBean, Initia builder.withoutJMXReporting(); } + if (sslEnabled) { + if (sslOptions == null) { + builder.withSSL(); + } else { + builder.withSSL(sslOptions); + } + } + cluster = builder.build(); + if (hostStateListener != null) { + cluster.register(hostStateListener); + } + + if (latencyTracker != null) { + cluster.register(latencyTracker); + } + generateSpecificationsFromFactoryBeans(); executeSpecsAndScripts(keyspaceCreations, startupScripts); @@ -373,4 +397,32 @@ public class CassandraClusterFactoryBean implements FactoryBean, Initia public void setJmxReportingEnabled(boolean jmxReportingEnabled) { this.jmxReportingEnabled = jmxReportingEnabled; } + + /** + * @param sslEnabled The sslEnabled to set. + */ + public void setSslEnabled(boolean sslEnabled) { + this.sslEnabled = sslEnabled; + } + + /** + * @param sslOptions The sslOptions to set. + */ + public void setSslOptions(SSLOptions sslOptions) { + this.sslOptions = sslOptions; + } + + /** + * @param hostStateListener The hostStateListener to set. + */ + public void setHostStateListener(Host.StateListener hostStateListener) { + this.hostStateListener = hostStateListener; + } + + /** + * @param latencyTracker The latencyTracker to set. + */ + public void setLatencyTracker(LatencyTracker latencyTracker) { + this.latencyTracker = latencyTracker; + } } diff --git a/spring-cassandra/src/main/java/org/springframework/cassandra/config/xml/CassandraClusterParser.java b/spring-cassandra/src/main/java/org/springframework/cassandra/config/xml/CassandraClusterParser.java index aafdd3a38..b750d4faf 100644 --- a/spring-cassandra/src/main/java/org/springframework/cassandra/config/xml/CassandraClusterParser.java +++ b/spring-cassandra/src/main/java/org/springframework/cassandra/config/xml/CassandraClusterParser.java @@ -110,11 +110,6 @@ public class CassandraClusterParser extends AbstractBeanDefinitionParser { builder.addPropertyValue("compressionType", compression); } - String authProvider = element.getAttribute("auth-info-provider-ref"); - if (StringUtils.hasText(authProvider)) { - builder.addPropertyReference("authProvider", authProvider); - } - String username = element.getAttribute("username"); if (StringUtils.hasText(username)) { builder.addPropertyValue("username", username); @@ -140,6 +135,16 @@ public class CassandraClusterParser extends AbstractBeanDefinitionParser { builder.addPropertyValue("jmxReportingEnabled", jmxReportingEnabled); } + String sslEnabled = element.getAttribute("sslEnabled"); + if (StringUtils.hasText(sslEnabled)) { + builder.addPropertyValue("sslEnabled", sslEnabled); + } + + String authProvider = element.getAttribute("auth-info-provider-ref"); + if (StringUtils.hasText(authProvider)) { + builder.addPropertyReference("authProvider", authProvider); + } + String loadBalancingPolicy = element.getAttribute("load-balancing-policy-ref"); if (StringUtils.hasText(loadBalancingPolicy)) { builder.addPropertyReference("loadBalancingPolicy", loadBalancingPolicy); @@ -155,6 +160,21 @@ public class CassandraClusterParser extends AbstractBeanDefinitionParser { builder.addPropertyReference("retryPolicy", retryPolicy); } + String sslOptions = element.getAttribute("ssl-options-ref"); + if (StringUtils.hasText(sslOptions)) { + builder.addPropertyReference("sslOptions", sslOptions); + } + + String hostStateListener = element.getAttribute("host-state-listener-ref"); + if (StringUtils.hasText(hostStateListener)) { + builder.addPropertyReference("hostStateListener", hostStateListener); + } + + String latencyTracker = element.getAttribute("latency-tracker-ref"); + if (StringUtils.hasText(latencyTracker)) { + builder.addPropertyReference("latencyTracker", latencyTracker); + } + parseChildElements(element, context, builder); } diff --git a/spring-cassandra/src/main/resources/org/springframework/cassandra/config/spring-cassandra-1.0.xsd b/spring-cassandra/src/main/resources/org/springframework/cassandra/config/spring-cassandra-1.0.xsd index 45bb9ddfc..6eed171ab 100644 --- a/spring-cassandra/src/main/resources/org/springframework/cassandra/config/spring-cassandra-1.0.xsd +++ b/spring-cassandra/src/main/resources/org/springframework/cassandra/config/spring-cassandra-1.0.xsd @@ -136,23 +136,6 @@ The protocol compression option. Default is "NONE". - - - - - - - - - - - - - - - - + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + @@ -479,7 +540,7 @@ Provides the ability to specify replication factors by data center. default="SimpleStrategy"> diff --git a/spring-cassandra/src/test/java/org/springframework/cassandra/test/integration/AbstractEmbeddedCassandraIntegrationTest.java b/spring-cassandra/src/test/java/org/springframework/cassandra/test/integration/AbstractEmbeddedCassandraIntegrationTest.java index 3315c021a..6cbf0563a 100644 --- a/spring-cassandra/src/test/java/org/springframework/cassandra/test/integration/AbstractEmbeddedCassandraIntegrationTest.java +++ b/spring-cassandra/src/test/java/org/springframework/cassandra/test/integration/AbstractEmbeddedCassandraIntegrationTest.java @@ -22,7 +22,7 @@ public class AbstractEmbeddedCassandraIntegrationTest { static Logger log = LoggerFactory.getLogger(AbstractEmbeddedCassandraIntegrationTest.class); - protected static final String CASSANDRA_CONFIG = "spring-cassandra.yaml"; + protected static String CASSANDRA_CONFIG = "spring-cassandra.yaml"; protected static final String CASSANDRA_HOST = "localhost"; protected static final int CASSANDRA_NATIVE_PORT = 9042; diff --git a/spring-cassandra/src/test/java/org/springframework/cassandra/test/integration/config/xml/TestHostStateListener.java b/spring-cassandra/src/test/java/org/springframework/cassandra/test/integration/config/xml/TestHostStateListener.java new file mode 100644 index 000000000..9ea86c48c --- /dev/null +++ b/spring-cassandra/src/test/java/org/springframework/cassandra/test/integration/config/xml/TestHostStateListener.java @@ -0,0 +1,52 @@ +/* + * Copyright 2011-2014 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.cassandra.test.integration.config.xml; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import com.datastax.driver.core.Host; +import com.datastax.driver.core.Host.StateListener; + +/** + * @author David Webb + * + */ +public class TestHostStateListener implements StateListener { + + private final static Logger log = LoggerFactory.getLogger(TestHostStateListener.class); + + @Override + public void onAdd(Host host) { + log.info("Host Added: " + host.getAddress()); + } + + @Override + public void onUp(Host host) { + log.info("Host Up: " + host.getAddress()); + } + + @Override + public void onDown(Host host) { + log.info("Host Down: " + host.getAddress()); + } + + @Override + public void onRemove(Host host) { + log.info("Host Removed: " + host.getAddress()); + } + +} diff --git a/spring-cassandra/src/test/java/org/springframework/cassandra/test/integration/config/xml/TestLatencyTracker.java b/spring-cassandra/src/test/java/org/springframework/cassandra/test/integration/config/xml/TestLatencyTracker.java new file mode 100644 index 000000000..221e408be --- /dev/null +++ b/spring-cassandra/src/test/java/org/springframework/cassandra/test/integration/config/xml/TestLatencyTracker.java @@ -0,0 +1,37 @@ +/* + * Copyright 2011-2014 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.cassandra.test.integration.config.xml; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import com.datastax.driver.core.Host; +import com.datastax.driver.core.LatencyTracker; + +/** + * @author David Webb + * + */ +public class TestLatencyTracker implements LatencyTracker { + + private final static Logger log = LoggerFactory.getLogger(TestLatencyTracker.class); + + @Override + public void update(Host host, long newLatencyNanos) { + log.info("Latency Tracker: " + host.getAddress() + ", " + newLatencyNanos + " nanoseconds."); + } + +} diff --git a/spring-cassandra/src/test/resources/org/springframework/cassandra/test/integration/config/xml/PropertyPlaceholderNamespaceCreatingXmlConfigTest-context.xml b/spring-cassandra/src/test/resources/org/springframework/cassandra/test/integration/config/xml/PropertyPlaceholderNamespaceCreatingXmlConfigTest-context.xml index 81fbb2471..317317470 100644 --- a/spring-cassandra/src/test/resources/org/springframework/cassandra/test/integration/config/xml/PropertyPlaceholderNamespaceCreatingXmlConfigTest-context.xml +++ b/spring-cassandra/src/test/resources/org/springframework/cassandra/test/integration/config/xml/PropertyPlaceholderNamespaceCreatingXmlConfigTest-context.xml @@ -21,6 +21,10 @@ + + + + + retry-policy-ref="retryPolicy" + host-state-listener-ref="hostStateListener" + latency-tracker-ref="latencyTracker">