diff --git a/spring-cql/src/main/java/org/springframework/cassandra/config/PoolingOptionsFactoryBean.java b/spring-cql/src/main/java/org/springframework/cassandra/config/PoolingOptionsFactoryBean.java index e332e92fe..55706dcaf 100644 --- a/spring-cql/src/main/java/org/springframework/cassandra/config/PoolingOptionsFactoryBean.java +++ b/spring-cql/src/main/java/org/springframework/cassandra/config/PoolingOptionsFactoryBean.java @@ -15,7 +15,8 @@ */ package org.springframework.cassandra.config; -import org.springframework.beans.factory.DisposableBean; +import java.util.concurrent.Executor; + import org.springframework.beans.factory.FactoryBean; import org.springframework.beans.factory.InitializingBean; @@ -23,218 +24,631 @@ import com.datastax.driver.core.HostDistance; import com.datastax.driver.core.PoolingOptions; /** - * Pooling Options Factory Bean. + * Spring {@link FactoryBean} for the Cassandra Java driver {@link PoolingOptions}. * * @author Matthew T. Adams * @author David Webb * @author Mark Paluch + * @author John Blum + * @see org.springframework.beans.factory.FactoryBean + * @see org.springframework.beans.factory.InitializingBean + * @see com.datastax.driver.core.PoolingOptions */ -public class PoolingOptionsFactoryBean implements FactoryBean, InitializingBean, DisposableBean { +@SuppressWarnings("unused") +public class PoolingOptionsFactoryBean implements FactoryBean, InitializingBean { - private Integer localMinSimultaneousRequests; - private Integer localMaxSimultaneousRequests; + private Executor initializationExecutor; + + private Integer heartbeatIntervalSeconds; + private Integer idleTimeoutSeconds; private Integer localCoreConnections; private Integer localMaxConnections; - private Integer remoteMinSimultaneousRequests; - private Integer remoteMaxSimultaneousRequests; + private Integer localMaxSimultaneousRequests; + private Integer localMinSimultaneousRequests; + private Integer poolTimeoutMilliseconds; private Integer remoteCoreConnections; private Integer remoteMaxConnections; + private Integer remoteMaxSimultaneousRequests; + private Integer remoteMinSimultaneousRequests; - PoolingOptions poolingOptions; - - @Override - public void destroy() throws Exception { - localMinSimultaneousRequests = null; - localMaxSimultaneousRequests = null; - localCoreConnections = null; - localMaxConnections = null; - remoteMinSimultaneousRequests = null; - remoteMaxSimultaneousRequests = null; - remoteCoreConnections = null; - remoteMaxConnections = null; - } + private PoolingOptions poolingOptions; + /* + * (non-Javadoc) + * @see org.springframework.beans.factory.InitializingBean#afterPropertiesSet() + */ @Override public void afterPropertiesSet() throws Exception { - poolingOptions = new PoolingOptions(); + poolingOptions = configureRemoteHostDistancePoolingOptions( + configureLocalHostDistancePoolingOptions(newPoolingOptions())); - if (localMaxConnections != null) { - poolingOptions.setMaxConnectionsPerHost(HostDistance.LOCAL, localMaxConnections); + if (heartbeatIntervalSeconds != null) { + poolingOptions.setHeartbeatIntervalSeconds(heartbeatIntervalSeconds); } - if (localCoreConnections != null) { - poolingOptions.setCoreConnectionsPerHost(HostDistance.LOCAL, localCoreConnections); + if (idleTimeoutSeconds != null) { + poolingOptions.setIdleTimeoutSeconds(idleTimeoutSeconds); } - if (localMinSimultaneousRequests != null) { - /* - * If the new min is greater than the current Max, set the current max to the new min first. - * This is enforced by the DSE Driver so you cannot set a new min/max together if either one falls outside of the default 25-100 range. - */ - int currentMax = poolingOptions.getNewConnectionThreshold(HostDistance.LOCAL); - if (currentMax < localMinSimultaneousRequests) { - poolingOptions.setNewConnectionThreshold(HostDistance.LOCAL, - localMinSimultaneousRequests); - } + if (initializationExecutor != null) { + poolingOptions.setInitializationExecutor(initializationExecutor); } - if (localMaxSimultaneousRequests != null) { - poolingOptions.setMaxRequestsPerConnection(HostDistance.LOCAL, localMaxSimultaneousRequests); + if (poolTimeoutMilliseconds != null) { + poolingOptions.setPoolTimeoutMillis(poolTimeoutMilliseconds); } - - if (remoteMaxConnections != null) { - poolingOptions.setMaxConnectionsPerHost(HostDistance.REMOTE, remoteMaxConnections); - } - - if (remoteCoreConnections != null) { - poolingOptions.setCoreConnectionsPerHost(HostDistance.REMOTE, remoteCoreConnections); - } - - if (remoteMinSimultaneousRequests != null) { - /* - * If the new min is greater than the current Max, set the current max to the new min first. - * This is enforced by the DSE Driver so you cannot set a new min/max together if either one falls outside of the default 25-100 range. - */ - int currentMax = poolingOptions.getNewConnectionThreshold(HostDistance.REMOTE); - if (currentMax < remoteMinSimultaneousRequests) { - poolingOptions.setNewConnectionThreshold(HostDistance.REMOTE, - remoteMinSimultaneousRequests); - } - } - - if (remoteMaxSimultaneousRequests != null) { - poolingOptions.setMaxRequestsPerConnection(HostDistance.REMOTE, - remoteMaxSimultaneousRequests); - } - } + /* + * (non-Javadoc) + * @see com.datastax.driver.core.PoolingOptions + */ + PoolingOptions newPoolingOptions() { + return new PoolingOptions(); + } + + /** + * Constructs and returns a {@link PoolingOptionsFactoryBean.HostDistancePoolingOptions} instance initialized + * with the {@link HostDistance#LOCAL}-based {@link PoolingOptions} as configured on this + * {@link PoolingOptionsFactoryBean}. + * + * @return {@link PoolingOptionsFactoryBean.HostDistancePoolingOptions} initialized with this + * {@link PoolingOptionsFactoryBean}'s {@link HostDistance#LOCAL}-based {@link PoolingOptions}. + * @see com.datastax.driver.core.HostDistance#LOCAL + * @see com.datastax.driver.core.PoolingOptions + * @see org.springframework.cassandra.config.PoolingOptionsFactoryBean.HostDistancePoolingOptions + * @see org.springframework.cassandra.config.PoolingOptionsFactoryBean.LocalHostDistancePoolingOptions + */ + protected HostDistancePoolingOptions newLocalHostDistancePoolingOptions() { + return LocalHostDistancePoolingOptions.create(getLocalCoreConnections(), getLocalMaxConnections(), + getLocalMaxSimultaneousRequests(), getLocalMinSimultaneousRequests()); + } + + /** + * Constructs and returns a {@link PoolingOptionsFactoryBean.HostDistancePoolingOptions} instance initialized + * with the {@link HostDistance#REMOTE}-based {@link PoolingOptions} as configured on this + * {@link PoolingOptionsFactoryBean}. + * + * @return {@link PoolingOptionsFactoryBean.HostDistancePoolingOptions} initialized with this + * {@link PoolingOptionsFactoryBean}'s {@link HostDistance#REMOTE}-based {@link PoolingOptions}. + * @see com.datastax.driver.core.HostDistance#REMOTE + * @see com.datastax.driver.core.PoolingOptions + * @see org.springframework.cassandra.config.PoolingOptionsFactoryBean.HostDistancePoolingOptions + * @see org.springframework.cassandra.config.PoolingOptionsFactoryBean.RemoteHostDistancePoolingOptions + */ + protected HostDistancePoolingOptions newRemoteHostDistancePoolingOptions() { + return RemoteHostDistancePoolingOptions.create(getRemoteCoreConnections(), getRemoteMaxConnections(), + getRemoteMaxSimultaneousRequests(), getRemoteMinSimultaneousRequests()); + } + + /** + * Configures the {@link HostDistance#LOCAL} connection settings on the given {@link PoolingOptions}. + * + * @param poolingOptions the {@link PoolingOptions} to configure. + * @return the given {@link PoolingOptions}. + * @see com.datastax.driver.core.HostDistance#LOCAL + * @see com.datastax.driver.core.PoolingOptions + * @see #newLocalHostDistancePoolingOptions() + */ + protected PoolingOptions configureLocalHostDistancePoolingOptions(PoolingOptions poolingOptions) { + return newLocalHostDistancePoolingOptions().configure(poolingOptions); + } + + /** + * Configures the {@link HostDistance#REMOTE} connection settings on the given {@link PoolingOptions}. + * + * @param poolingOptions the {@link PoolingOptions} to configure. + * @return the given {@link PoolingOptions}. + * @see com.datastax.driver.core.HostDistance#REMOTE + * @see com.datastax.driver.core.PoolingOptions + * @see #newRemoteHostDistancePoolingOptions() + */ + protected PoolingOptions configureRemoteHostDistancePoolingOptions(PoolingOptions poolingOptions) { + return newRemoteHostDistancePoolingOptions().configure(poolingOptions); + } + + /* + * (non-Javadoc) + * @see org.springframework.beans.factory.FactoryBean#getObject() + */ @Override public PoolingOptions getObject() throws Exception { return poolingOptions; } + /* + * (non-Javadoc) + * @see org.springframework.beans.factory.FactoryBean#getObjectType() + */ @Override public Class getObjectType() { - return PoolingOptions.class; + return (poolingOptions != null ? poolingOptions.getClass() : PoolingOptions.class); } + /* + * (non-Javadoc) + * @see org.springframework.beans.factory.FactoryBean#isSingleton() + */ @Override public boolean isSingleton() { return true; } /** - * @return Returns the localMinSimultaneousRequests. + * Sets the heart beat interval, after which a message is sent on an idle connection to make sure it's still alive. + * + * @param heartbeatIntervalSeconds interval in seconds between heartbeat messages to keep idle connections alive. */ - public Integer getLocalMinSimultaneousRequests() { - return localMinSimultaneousRequests; + public void setHeartbeatIntervalSeconds(Integer heartbeatIntervalSeconds) { + this.heartbeatIntervalSeconds = heartbeatIntervalSeconds; } /** - * @param localMinSimultaneousRequests The localMinSimultaneousRequests to set. + * Gets the heart beat interval, after which a message is sent on an idle connection to make sure it's still alive. + * + * @return the {@code heartbeatIntervalSeconds}. */ - public void setLocalMinSimultaneousRequests(Integer localMinSimultaneousRequests) { - this.localMinSimultaneousRequests = localMinSimultaneousRequests; + public Integer getHeartbeatIntervalSeconds() { + return heartbeatIntervalSeconds; } /** - * @return Returns the localMaxSimultaneousRequests. + * Sets the timeout before an idle connection is removed. + * + * @param idleTimeoutSeconds idle timeout in seconds before a connection is removed. */ - public Integer getLocalMaxSimultaneousRequests() { - return localMaxSimultaneousRequests; + public void setIdleTimeoutSeconds(Integer idleTimeoutSeconds) { + this.idleTimeoutSeconds = idleTimeoutSeconds; } /** - * @param localMaxSimultaneousRequests The localMaxSimultaneousRequests to set. + * Get the timeout before an idle connection is removed. + * + * @return the {@code idleTimeoutSeconds}. */ - public void setLocalMaxSimultaneousRequests(Integer localMaxSimultaneousRequests) { - this.localMaxSimultaneousRequests = localMaxSimultaneousRequests; + public Integer getIdleTimeoutSeconds() { + return idleTimeoutSeconds; } /** - * @return Returns the localCoreConnections. + * Sets the {@link Executor} to use for connection initialization. + * + * @param initializationExecutor {@link Executor} used to initialize the connection. */ - public Integer getLocalCoreConnections() { - return localCoreConnections; + public void setInitializationExecutor(Executor initializationExecutor) { + this.initializationExecutor = initializationExecutor; } /** - * @param localCoreConnections The localCoreConnections to set. + * Gets the {@link Executor} to use for connection initialization. + * + * @return the {@code initializationExecutor}. + */ + public Executor getInitializationExecutor() { + return initializationExecutor; + } + + /** + * Sets the timeout when trying to acquire a connection from a host's pool. + * + * @param poolTimeoutMilliseconds timeout in milliseconds used to acquire a connection from the host's pool. + */ + public void setPoolTimeoutMilliseconds(Integer poolTimeoutMilliseconds) { + this.poolTimeoutMilliseconds = poolTimeoutMilliseconds; + } + + /** + * Gets the timeout when trying to acquire a connection from a host's pool. + * + * @return the {@code poolTimeoutMilliseconds}. + */ + public Integer getPoolTimeoutMilliseconds() { + return poolTimeoutMilliseconds; + } + + /** + * Sets the core number of connections per host for the {@link HostDistance#LOCAL} scope. + * + * @param localCoreConnections core number of local connections per host. */ public void setLocalCoreConnections(Integer localCoreConnections) { this.localCoreConnections = localCoreConnections; } /** - * @return Returns the localMaxConnections. + * Gets the core number of connections per host for the {@link HostDistance#LOCAL} scope. + * + * @return the {@code localCoreConnections). */ - public Integer getLocalMaxConnections() { - return localMaxConnections; + public Integer getLocalCoreConnections() { + return localCoreConnections; } /** - * @param localMaxConnections The localMaxConnections to set. + * Sets the maximum number of connections per host for the {@link HostDistance#LOCAL} scope. + * + * @param localMaxConnections max number of local connections per host. */ public void setLocalMaxConnections(Integer localMaxConnections) { this.localMaxConnections = localMaxConnections; } /** - * @return Returns the remoteMinSimultaneousRequests. + * Gets the maximum number of connections per host for the {@link HostDistance#LOCAL} scope. + * + * @return the {@code localMaxConnections}. */ - public Integer getRemoteMinSimultaneousRequests() { - return remoteMinSimultaneousRequests; + public Integer getLocalMaxConnections() { + return localMaxConnections; } /** - * @param remoteMinSimultaneousRequests The remoteMinSimultaneousRequests to set. + * Sets the maximum number of requests per connection for the {@link HostDistance#LOCAL} scope. + * + * @param localMaxSimultaneousRequests max number of requests for local connections. */ - public void setRemoteMinSimultaneousRequests(Integer remoteMinSimultaneousRequests) { - this.remoteMinSimultaneousRequests = remoteMinSimultaneousRequests; + public void setLocalMaxSimultaneousRequests(Integer localMaxSimultaneousRequests) { + this.localMaxSimultaneousRequests = localMaxSimultaneousRequests; } /** - * @return Returns the remoteMaxSimultaneousRequests. + * Gets the maximum number of requests per connection for the {@link HostDistance#LOCAL} scope. + * + * @return the {@code localMaxSimultaneousRequests}. */ - public Integer getRemoteMaxSimultaneousRequests() { - return remoteMaxSimultaneousRequests; + public Integer getLocalMaxSimultaneousRequests() { + return localMaxSimultaneousRequests; } /** - * @param remoteMaxSimultaneousRequests The remoteMaxSimultaneousRequests to set. + * Sets the threshold that triggers the creation of a new connection to a host + * for the {@link HostDistance#LOCAL} scope. + * + * @param localMinSimultaneousRequests threshold triggering the creation of local connections to a host. */ - public void setRemoteMaxSimultaneousRequests(Integer remoteMaxSimultaneousRequests) { - this.remoteMaxSimultaneousRequests = remoteMaxSimultaneousRequests; + public void setLocalMinSimultaneousRequests(Integer localMinSimultaneousRequests) { + this.localMinSimultaneousRequests = localMinSimultaneousRequests; } /** - * @return Returns the remoteCoreConnections. + * Gets the threshold that triggers the creation of a new connection to a host + * for the {@link HostDistance#LOCAL} scope. + * + * @return the {@code localMinSimultaneousRequests}. */ - public Integer getRemoteCoreConnections() { - return remoteCoreConnections; + public Integer getLocalMinSimultaneousRequests() { + return localMinSimultaneousRequests; } /** - * @param remoteCoreConnections The remoteCoreConnections to set. + * Sets the core number of connections per host for the {@link HostDistance#REMOTE} scope. + * + * @param remoteCoreConnections core number of remote connections per host. */ public void setRemoteCoreConnections(Integer remoteCoreConnections) { this.remoteCoreConnections = remoteCoreConnections; } /** - * @return Returns the remoteMaxConnections. + * Gets the core number of connections per host for the {@link HostDistance#REMOTE} scope. + * + * @return the {@code remoteCoreConnections). + */ + public Integer getRemoteCoreConnections() { + return remoteCoreConnections; + } + + /** + * Sets the maximum number of connections per host for the {@link HostDistance#REMOTE} scope. + * + * @param remoteMaxConnections max number of remote connections per host. + */ + public void setRemoteMaxConnections(Integer remoteMaxConnections) { + this.remoteMaxConnections = remoteMaxConnections; + } + + /** + * Gets the maximum number of connections per host for the {@link HostDistance#REMOTE} scope. + * + * @return the {@code remoteMaxConnections}. */ public Integer getRemoteMaxConnections() { return remoteMaxConnections; } /** - * @param remoteMaxConnections The remoteMaxConnections to set. + * Sets the maximum number of requests per connection for the {@link HostDistance#REMOTE} scope. + * + * @param remoteMaxSimultaneousRequests max number of requests for local connections. */ - public void setRemoteMaxConnections(Integer remoteMaxConnections) { - this.remoteMaxConnections = remoteMaxConnections; + public void setRemoteMaxSimultaneousRequests(Integer remoteMaxSimultaneousRequests) { + this.remoteMaxSimultaneousRequests = remoteMaxSimultaneousRequests; } + /** + * Gets the maximum number of requests per connection for the {@link HostDistance#REMOTE} scope. + * + * @return the {@code remoteMaxSimultaneousRequests}. + */ + public Integer getRemoteMaxSimultaneousRequests() { + return remoteMaxSimultaneousRequests; + } + + /** + * Sets the threshold that triggers the creation of a new connection to a host + * for the {@link HostDistance#REMOTE} scope. + * + * @param remoteMinSimultaneousRequests threshold triggering the creation of remote connections to a host. + */ + public void setRemoteMinSimultaneousRequests(Integer remoteMinSimultaneousRequests) { + this.remoteMinSimultaneousRequests = remoteMinSimultaneousRequests; + } + + /** + * Gets the threshold that triggers the creation of a new connection to a host + * for the {@link HostDistance#REMOTE} scope. + * + * @return the {@code remoteMinSimultaneousRequests}. + */ + public Integer getRemoteMinSimultaneousRequests() { + return remoteMinSimultaneousRequests; + } + + /** + * The HostDistancePoolingOptions class models the {@link PoolingOptions} state and connection settings for a + * particular {@link HostDistance}. + * + * @see com.datastax.driver.core.HostDistance + * @see com.datastax.driver.core.PoolingOptions + */ + protected static abstract class HostDistancePoolingOptions { + + private final Integer coreConnectionsPerHost; + private final Integer maxConnectionsPerHost; + private final Integer maxRequestsPerConnection; + private final Integer newConnectionThreshold; + + /** + * Constructs an instance of {@link HostDistancePoolingOptions} with {@link PoolingOptions} connection settings + * specific to a particular {@link HostDistance}. + * + * @param coreConnectionsPerHost core number of connections per host. + * @param maxConnectionsPerHost maximum number of connections per host. + * @param maxRequestsPerConnection maximum number of requests per connection. + * @param newConnectionThreshold threshold that triggers the creation of a new connection to a host. + */ + protected HostDistancePoolingOptions(Integer coreConnectionsPerHost, Integer maxConnectionsPerHost, + Integer maxRequestsPerConnection, Integer newConnectionThreshold) { + + this.coreConnectionsPerHost = coreConnectionsPerHost; + this.maxConnectionsPerHost = maxConnectionsPerHost; + this.maxRequestsPerConnection = maxRequestsPerConnection; + this.newConnectionThreshold = newConnectionThreshold; + } + + /** + * Returns the {@link HostDistance} used to configure the specific {@link PoolingOptions} connection settings. + * + * @return a {@link HostDistance} used to configure the specific {@link PoolingOptions} connection settings. + * @see com.datastax.driver.core.HostDistance + */ + protected abstract HostDistance getHostDistance(); + + /* + * (non-Javadoc) + * @see com.datastax.driver.core.PoolingOptions#setCoreConnectionsPerHost(HostDistance, int) + */ + PoolingOptions setCoreConnectionsPerHost(PoolingOptions poolingOptions) { + if (coreConnectionsPerHost != null) { + poolingOptions.setCoreConnectionsPerHost(getHostDistance(), coreConnectionsPerHost); + } + + return poolingOptions; + } + + /** + * Returns the core number of connections per host. + * + * @return {@code coreConnectionsPerHost}. + * @see com.datastax.driver.core.PoolingOptions#getCoreConnectionsPerHost(HostDistance) + * @see #getHostDistance() + */ + protected Integer getCoreConnectionsPerHost() { + return coreConnectionsPerHost; + } + + /* + * (non-Javadoc) + * @see com.datastax.driver.core.PoolingOptions#setMaxConnectionsPerHost(HostDistance, int) + */ + PoolingOptions setMaxConnectionsPerHost(PoolingOptions poolingOptions) { + if (maxConnectionsPerHost != null) { + poolingOptions.setMaxConnectionsPerHost(getHostDistance(), maxConnectionsPerHost); + } + + return poolingOptions; + } + + /** + * Returns the maximum number of connections per host. + * + * @return {@code maxConnectionsPerHost}. + * @see com.datastax.driver.core.PoolingOptions#getMaxConnectionsPerHost(HostDistance) + * @see #getHostDistance() + */ + protected Integer getMaxConnectionsPerHost() { + return maxConnectionsPerHost; + } + + /* + * (non-Javadoc) + * @see com.datastax.driver.core.PoolingOptions#setMaxRequestsPerConnection(HostDistance, int) + */ + PoolingOptions setMaxRequestsPerConnection(PoolingOptions poolingOptions) { + if (maxRequestsPerConnection != null) { + poolingOptions.setMaxRequestsPerConnection(getHostDistance(), maxRequestsPerConnection); + } + + return poolingOptions; + } + + /** + * Returns the maximum number of requests per connection. + * + * @return {@code maxRequestsPerConnection}. + * @see com.datastax.driver.core.PoolingOptions#getMaxRequestsPerConnection(HostDistance) + * @see #getHostDistance() + */ + protected Integer getMaxRequestsPerConnection() { + return maxRequestsPerConnection; + } + + /* + * If the new min is greater than the current max, set the current max to the new min first. + * This is enforced by the DSE Driver so you cannot set a new min/max together if either one falls outside + * of the default 25-100 range. + * + * @see com.datastax.driver.core.PoolingOptions#setNewConnectionThreshold(HostDistance, int) + */ + PoolingOptions setNewConnectionThreshold(PoolingOptions poolingOptions) { + if (newConnectionThreshold != null) { + int currentNewConnectionThreshold = poolingOptions.getNewConnectionThreshold(getHostDistance()); + + if (currentNewConnectionThreshold < newConnectionThreshold) { + poolingOptions.setNewConnectionThreshold(getHostDistance(), newConnectionThreshold); + } + } + + return poolingOptions; + } + + /** + * Returns the threshold that triggers the creation of a new connection to a host. + * + * @return {@code newConnectionThreshold}. + * @see com.datastax.driver.core.PoolingOptions#getNewConnectionThreshold(HostDistance) + * @see #getHostDistance() + */ + protected Integer getNewConnectionThreshold() { + return newConnectionThreshold; + } + + /* + * (non-Javadoc) + * @see com.datastax.driver.core.PoolingOptions + */ + PoolingOptions configure(PoolingOptions poolingOptions) { + + // order is important here; max properties must be set first + setMaxConnectionsPerHost(poolingOptions); + setCoreConnectionsPerHost(poolingOptions); + setMaxRequestsPerConnection(poolingOptions); + setNewConnectionThreshold(poolingOptions); + + return poolingOptions; + } + } + + /* + * (non-Javadoc) + * @see HostDistancePoolingOptions + * @see com.datastax.driver.core.PoolingOptions + * @see com.datastax.driver.core.HostDistance#LOCAL + */ + static class LocalHostDistancePoolingOptions extends HostDistancePoolingOptions { + + /** + * Creates an instance of {@link LocalHostDistancePoolingOptions} initialized with {@link PoolingOptions} + * based on {@link HostDistance#LOCAL}. + * + * @param coreConnectionsPerHost core number of connections per host. + * @param maxConnectionsPerHost maximum number of connections per host. + * @param maxRequestsPerConnection maximum number of requests per connection. + * @param newConnectionThreshold threshold that triggers the creation of a new connection to a host. + */ + static LocalHostDistancePoolingOptions create(Integer coreConnectionsPerHost, Integer maxConnectionsPerHost, + Integer maxRequestsPerConnection, Integer newConnectionThreshold) { + + return new LocalHostDistancePoolingOptions(coreConnectionsPerHost, maxConnectionsPerHost, + maxRequestsPerConnection, newConnectionThreshold); + } + + /** + * Constructs an instance of {@link LocalHostDistancePoolingOptions} initialized with {@link PoolingOptions} + * based on {@link HostDistance#LOCAL}. + * + * @param coreConnectionsPerHost core number of connections per host. + * @param maxConnectionsPerHost maximum number of connections per host. + * @param maxRequestsPerConnection maximum number of requests per connection. + * @param newConnectionThreshold threshold that triggers the creation of a new connection to a host. + */ + LocalHostDistancePoolingOptions(Integer coreConnectionsPerHost, Integer maxConnectionsPerHost, + Integer maxRequestsPerConnection, Integer newConnectionThreshold) { + + super(coreConnectionsPerHost, maxConnectionsPerHost, maxRequestsPerConnection, newConnectionThreshold); + } + + /** + * Returns {@link HostDistance#LOCAL} to configure local-based {@link PoolingOptions} connection settings. + * + * @return {@link HostDistance#LOCAL} to configure local-based {@link PoolingOptions} connection settings. + * @see com.datastax.driver.core.HostDistance#LOCAL + */ + @Override + protected HostDistance getHostDistance() { + return HostDistance.LOCAL; + } + } + + /* + * (non-Javadoc) + * @see HostDistancePoolingOptions + * @see com.datastax.driver.core.PoolingOptions + * @see com.datastax.driver.core.HostDistance#REMOTE + */ + static class RemoteHostDistancePoolingOptions extends HostDistancePoolingOptions { + + /** + * Creates an instance of {@link RemoteHostDistancePoolingOptions} initialized with {@link PoolingOptions} + * based on {@link HostDistance#REMOTE}. + * + * @param coreConnectionsPerHost core number of connections per host. + * @param maxConnectionsPerHost maximum number of connections per host. + * @param maxRequestsPerConnection maximum number of requests per connection. + * @param newConnectionThreshold threshold that triggers the creation of a new connection to a host. + */ + static RemoteHostDistancePoolingOptions create(Integer coreConnectionsPerHost, Integer maxConnectionsPerHost, + Integer maxRequestsPerConnection, Integer newConnectionThreshold) { + + return new RemoteHostDistancePoolingOptions(coreConnectionsPerHost, maxConnectionsPerHost, + maxRequestsPerConnection, newConnectionThreshold); + } + + /** + * Constructs an instance of {@link RemoteHostDistancePoolingOptions} initialized with {@link PoolingOptions} + * based on {@link HostDistance#REMOTE}. + * + * @param coreConnectionsPerHost core number of connections per host. + * @param maxConnectionsPerHost maximum number of connections per host. + * @param maxRequestsPerConnection maximum number of requests per connection. + * @param newConnectionThreshold threshold that triggers the creation of a new connection to a host. + */ + RemoteHostDistancePoolingOptions(Integer coreConnectionsPerHost, Integer maxConnectionsPerHost, + Integer maxRequestsPerConnection, Integer newConnectionThreshold) { + + super(coreConnectionsPerHost, maxConnectionsPerHost, maxRequestsPerConnection, newConnectionThreshold); + } + + /** + * Returns {@link HostDistance#REMOTE} to configure remote-based {@link PoolingOptions} connection settings. + * + * @return {@link HostDistance#REMOTE} to configure remote-based {@link PoolingOptions} connection settings. + * @see com.datastax.driver.core.HostDistance#REMOTE + */ + @Override + protected HostDistance getHostDistance() { + return HostDistance.REMOTE; + } + } } diff --git a/spring-cql/src/main/java/org/springframework/cassandra/config/xml/CassandraCqlClusterParser.java b/spring-cql/src/main/java/org/springframework/cassandra/config/xml/CassandraCqlClusterParser.java index 2b05e33dd..fa4d07c1c 100644 --- a/spring-cql/src/main/java/org/springframework/cassandra/config/xml/CassandraCqlClusterParser.java +++ b/spring-cql/src/main/java/org/springframework/cassandra/config/xml/CassandraCqlClusterParser.java @@ -1,12 +1,12 @@ /* * Copyright 2013-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. @@ -15,6 +15,8 @@ */ package org.springframework.cassandra.config.xml; +import static org.springframework.cassandra.config.xml.ParsingUtils.*; + import java.util.ArrayList; import java.util.List; @@ -37,20 +39,15 @@ import org.springframework.util.StringUtils; import org.springframework.util.xml.DomUtils; import org.w3c.dom.Element; -import com.datastax.driver.core.HostDistance; import com.datastax.driver.core.PoolingOptions; import com.datastax.driver.core.SocketOptions; -import static org.springframework.cassandra.config.xml.ParsingUtils.addOptionalPropertyReference; -import static org.springframework.cassandra.config.xml.ParsingUtils.addOptionalPropertyValue; -import static org.springframework.cassandra.config.xml.ParsingUtils.addRequiredPropertyValue; -import static org.springframework.cassandra.config.xml.ParsingUtils.getSourceBeanDefinition; - /** * Parses the {@literal } element of the XML Configuration. - * + * * @author Matthew T. Adams * @author David Webb + * @author John Blum */ public class CassandraCqlClusterParser extends AbstractBeanDefinitionParser { @@ -59,23 +56,23 @@ public class CassandraCqlClusterParser extends AbstractBeanDefinitionParser { throws BeanDefinitionStoreException { String id = super.resolveId(element, definition, parserContext); - return StringUtils.hasText(id) ? id : DefaultCqlBeanNames.CLUSTER; + + return (StringUtils.hasText(id) ? id : DefaultCqlBeanNames.CLUSTER); } @Override protected AbstractBeanDefinition parseInternal(Element element, ParserContext parserContext) { - BeanDefinitionBuilder builder = BeanDefinitionBuilder.genericBeanDefinition(CassandraCqlClusterFactoryBean.class); - builder.getRawBeanDefinition().setSource(parserContext.extractSource(element)); - builder.getRawBeanDefinition().setDestroyMethodName("destroy"); - if (parserContext.isNested()) { - // Inner bean definition must receive same scope as containing bean. - builder.setScope(parserContext.getContainingBeanDefinition().getScope()); - } + BeanDefinitionBuilder builder = BeanDefinitionBuilder.genericBeanDefinition( + CassandraCqlClusterFactoryBean.class); - if (parserContext.isDefaultLazyInit()) { - // Default-lazy-init applies to custom bean definitions as well. - builder.setLazyInit(true); + builder.setLazyInit(parserContext.isDefaultLazyInit()); + builder.getRawBeanDefinition().setDestroyMethodName("destroy"); + builder.getRawBeanDefinition().setSource(parserContext.extractSource(element)); + + if (parserContext.isNested()) { + // inner bean definitions must have same scope as containing bean + builder.setScope(parserContext.getContainingBeanDefinition().getScope()); } doParse(element, parserContext, builder); @@ -84,75 +81,72 @@ public class CassandraCqlClusterParser extends AbstractBeanDefinitionParser { } /** - * Parses the attributes on the top level element, then parses all children. - * - * @param element The Element being parsed - * @param context The Parser Context - * @param builder The parent {@link BeanDefinitionBuilder} + * Parses cluster meta-data. + * + * @param element {@link Element} to parse. + * @param parserContext XML parser context and state. + * @param builder parent {@link BeanDefinitionBuilder}. */ - protected void doParse(Element element, ParserContext context, BeanDefinitionBuilder builder) { + protected void doParse(Element element, ParserContext parserContext, BeanDefinitionBuilder builder) { - addOptionalPropertyValue(builder, "contactPoints", element, "contact-points", null); - addOptionalPropertyValue(builder, "port", element, "port", null); - addOptionalPropertyValue(builder, "compressionType", element, "compression", null); - addOptionalPropertyValue(builder, "username", element, "username", null); - addOptionalPropertyValue(builder, "password", element, "password", null); - addOptionalPropertyValue(builder, "metricsEnabled", element, "metrics-enabled", null); - addOptionalPropertyValue(builder, "jmxReportingEnabled", element, "jmx-reporting-enabled", null); - addOptionalPropertyValue(builder, "sslEnabled", element, "ssl-enabled", null); + addOptionalPropertyReference(builder, "authProvider", element, "auth-info-provider-ref"); + addOptionalPropertyReference(builder, "hostStateListener", element, "host-state-listener-ref"); + addOptionalPropertyReference(builder, "latencyTracker", element, "latency-tracker-ref"); + addOptionalPropertyReference(builder, "loadBalancingPolicy", element, "load-balancing-policy-ref"); + addOptionalPropertyReference(builder, "reconnectionPolicy", element, "reconnection-policy-ref"); + addOptionalPropertyReference(builder, "retryPolicy", element, "retry-policy-ref"); + addOptionalPropertyReference(builder, "sslOptions", element, "ssl-options-ref"); - addOptionalPropertyReference(builder, "authProvider", element, "auth-info-provider-ref", null); - addOptionalPropertyReference(builder, "loadBalancingPolicy", element, "load-balancing-policy-ref", null); - addOptionalPropertyReference(builder, "reconnectionPolicy", element, "reconnection-policy-ref", null); - addOptionalPropertyReference(builder, "retryPolicy", element, "retry-policy-ref", null); - addOptionalPropertyReference(builder, "sslOptions", element, "ssl-options-ref", null); - addOptionalPropertyReference(builder, "hostStateListener", element, "host-state-listener-ref", null); - addOptionalPropertyReference(builder, "latencyTracker", element, "latency-tracker-ref", null); + addOptionalPropertyValue(builder, "contactPoints", element, "contact-points"); + addOptionalPropertyValue(builder, "compressionType", element, "compression"); + addOptionalPropertyValue(builder, "jmxReportingEnabled", element, "jmx-reporting-enabled"); + addOptionalPropertyValue(builder, "metricsEnabled", element, "metrics-enabled"); + addOptionalPropertyValue(builder, "password", element, "password"); + addOptionalPropertyValue(builder, "port", element, "port"); + addOptionalPropertyValue(builder, "sslEnabled", element, "ssl-enabled"); + addOptionalPropertyValue(builder, "username", element, "username"); - parseChildElements(element, context, builder); + parseChildElements(element, parserContext, builder); } /** - * Parse the Child Element of {@link DefaultCqlBeanNames.CLUSTER} - * - * @param element The Element being parsed - * @param context The Parser Context - * @param builder The parent {@link BeanDefinitionBuilder} + * Parses child elements of cluster. + * + * @param element {@link Element} to parse. + * @param parserContext XML parser context and state. + * @param builder parent {@link BeanDefinitionBuilder}. */ - protected void parseChildElements(Element element, ParserContext context, BeanDefinitionBuilder builder) { + protected void parseChildElements(Element element, ParserContext parserContext, BeanDefinitionBuilder builder) { ManagedSet keyspaceActionSpecificationBeanDefinitions = new ManagedSet(); + List startupScripts = new ArrayList(); List shutdownScripts = new ArrayList(); - List elements = DomUtils.getChildElements(element); - BeanDefinition keyspaceActionSpecificationBeanDefinition = null; + BeanDefinitionBuilder poolingOptionsBuilder = BeanDefinitionBuilder.genericBeanDefinition( + PoolingOptionsFactoryBean.class); - /* - * PoolingOptionsBuilder has two potential parsing cycles so it is defined - * before the child elements are iterated over, then converted to a BeanDefinition - * just in time. - */ - BeanDefinitionBuilder poolingOptionsBuilder = null; + addOptionalPropertyReference(poolingOptionsBuilder, "initializationExecutor", + element, "initialization-executor-ref"); - /* - * Parse each of the child elements - */ - for (Element subElement : elements) { + addOptionalPropertyValue(poolingOptionsBuilder, "heartbeatIntervalSeconds", element, "heartbeat-interval-seconds"); + addOptionalPropertyValue(poolingOptionsBuilder, "idleTimeoutSeconds", element, "idle-timeout-seconds"); + addOptionalPropertyValue(poolingOptionsBuilder, "poolTimeoutMilliseconds", element, "pool-timeout-milliseconds"); + + // parse child elements + for (Element subElement : DomUtils.getChildElements(element)) { String name = subElement.getLocalName(); - if ("local-pooling-options".equals(name)) { - poolingOptionsBuilder = parsePoolingOptions(subElement, poolingOptionsBuilder, HostDistance.LOCAL); + if ("keyspace".equals(name)) { + keyspaceActionSpecificationBeanDefinitions.add( + newKeyspaceActionSpecificationBeanDefinition(subElement, parserContext)); + } else if ("local-pooling-options".equals(name)) { + parseLocalPoolingOptions(subElement, poolingOptionsBuilder); } else if ("remote-pooling-options".equals(name)) { - poolingOptionsBuilder = parsePoolingOptions(subElement, poolingOptionsBuilder, HostDistance.REMOTE); + parseRemotePoolingOptions(subElement, poolingOptionsBuilder); } else if ("socket-options".equals(name)) { - builder.addPropertyValue("socketOptions", getSocketOptionsBeanDefinition(subElement, context)); - } else if ("keyspace".equals(name)) { - - keyspaceActionSpecificationBeanDefinition = getKeyspaceSpecificationBeanDefinition(subElement, context); - keyspaceActionSpecificationBeanDefinitions.add(keyspaceActionSpecificationBeanDefinition); - + builder.addPropertyValue("socketOptions", newSocketOptionsBeanDefinition(subElement, parserContext)); } else if ("startup-cql".equals(name)) { startupScripts.add(parseScript(subElement)); } else if ("shutdown-cql".equals(name)) { @@ -160,59 +154,64 @@ public class CassandraCqlClusterParser extends AbstractBeanDefinitionParser { } } - /* - * If the PoolingOptionsBuilder was initialized during parsing, process it now. - */ - if (poolingOptionsBuilder != null) { - builder.addPropertyValue("poolingOptions", getSourceBeanDefinition(poolingOptionsBuilder, context, element)); - } + builder.addPropertyValue("keyspaceSpecifications", newKeyspaceSetFlattenerBeanDefinition( + element, parserContext, keyspaceActionSpecificationBeanDefinitions)); + + builder.addPropertyValue("poolingOptions", getSourceBeanDefinition( + poolingOptionsBuilder, parserContext, element)); - builder.addPropertyValue("keyspaceSpecifications", - getKeyspaceSetFlattenerBeanDefinition(element, context, keyspaceActionSpecificationBeanDefinitions)); builder.addPropertyValue("startupScripts", startupScripts); builder.addPropertyValue("shutdownScripts", startupScripts); } /** - * Create the Single Factory Bean that will flatten all Set> - * - * @param element The Element being parsed - * @param context The Parser Context - * @param keyspaceActionSpecificationBeanDefinitions The List of Definitions to flatten - * @return A single level List of KeyspaceActionSpecifications + * Returns a {@link BeanDefinition} for a {@link KeyspaceActionSpecification} object. + * + * @param element Element being parsed. + * @param parserContext XML parser context and state. + * @return the {@link BeanDefinition} or {@literal null} if action is not given. */ - private Object getKeyspaceSetFlattenerBeanDefinition(Element element, ParserContext context, - ManagedSet keyspaceActionSpecificationBeanDefinitions) { + BeanDefinition newKeyspaceActionSpecificationBeanDefinition(Element element, ParserContext parserContext) { - BeanDefinitionBuilder flat = BeanDefinitionBuilder.genericBeanDefinition(MultiLevelSetFlattenerFactoryBean.class); - flat.addPropertyValue("multiLevelSet", keyspaceActionSpecificationBeanDefinitions); - return getSourceBeanDefinition(flat, context, element); + BeanDefinitionBuilder builder = BeanDefinitionBuilder.genericBeanDefinition( + KeyspaceActionSpecificationFactoryBean.class); + // add required replication defaults + addRequiredPropertyValue(builder, "replicationStrategy", + KeyspaceAttributes.DEFAULT_REPLICATION_STRATEGY.name()); + + addRequiredPropertyValue(builder, "replicationFactor", + String.valueOf(KeyspaceAttributes.DEFAULT_REPLICATION_FACTOR)); + + addRequiredPropertyValue(builder, "name", element, "name"); + addOptionalPropertyValue(builder, "durableWrites", element, "durable-writes", "false"); + addRequiredPropertyValue(builder, "action", element, "action"); + + parseReplication(DomUtils.getChildElementByTagName(element, "replication"), builder); + + return getSourceBeanDefinition(builder, parserContext, element); } /** * Parses the keyspace replication options and adds them to the supplied {@link BeanDefinitionBuilder}. - * - * @param element The Element being parsed + * + * @param element {@link Element} to parse. * @param builder The {@link BeanDefinitionBuilder} to add the replication to */ - protected void parseReplication(Element element, BeanDefinitionBuilder builder) { + void parseReplication(Element element, BeanDefinitionBuilder builder) { ManagedList networkTopologyDataCenters = new ManagedList(); ManagedList networkTopologyReplicationFactors = new ManagedList(); if (element != null) { - addOptionalPropertyValue(builder, "replicationStrategy", element, "class", - KeyspaceAttributes.DEFAULT_REPLICATION_STRATEGY.name()); - addOptionalPropertyValue(builder, "replicationFactor", element, "replication-factor", "" - + KeyspaceAttributes.DEFAULT_REPLICATION_FACTOR); + KeyspaceAttributes.DEFAULT_REPLICATION_STRATEGY.name()); - /* - * DataCenters only apply to NetworkTolopogyStrategy - */ - List dcElements = DomUtils.getChildElementsByTagName(element, "data-center"); - for (Element dataCenter : dcElements) { + addOptionalPropertyValue(builder, "replicationFactor", element, "replication-factor", + String.valueOf(KeyspaceAttributes.DEFAULT_REPLICATION_FACTOR)); + + // DataCenters only apply to NetworkTopologyStrategy + for (Element dataCenter : DomUtils.getChildElementsByTagName(element, "data-center")) { networkTopologyDataCenters.add(dataCenter.getAttribute("name")); networkTopologyReplicationFactors.add(dataCenter.getAttribute("replication-factor")); } @@ -223,94 +222,82 @@ public class CassandraCqlClusterParser extends AbstractBeanDefinitionParser { } /** - * Parse CQL Script Elements - * - * @param element The Element being parsed - * @return + * Create the Single Factory Bean that will flatten all Set> + * + * @param element {@link Element} to parse. + * @param parserContext XML parser context and state. + * @param keyspaceActionSpecificationBeanDefinitions The List of Definitions to flatten + * @return A single level List of KeyspaceActionSpecifications */ - protected String parseScript(Element element) { + Object newKeyspaceSetFlattenerBeanDefinition(Element element, ParserContext parserContext, + ManagedSet keyspaceActionSpecificationBeanDefinitions) { + + BeanDefinitionBuilder builder = BeanDefinitionBuilder.genericBeanDefinition( + MultiLevelSetFlattenerFactoryBean.class); + + builder.addPropertyValue("multiLevelSet", keyspaceActionSpecificationBeanDefinitions); + + return getSourceBeanDefinition(builder, parserContext, element); + } + + /** + * Parses local pooling options. + * + * @param element {@link Element} to parse. + * @param builder {@link BeanDefinitionBuilder} used to build a {@link PoolingOptions} {@link BeanDefinition}. + */ + void parseLocalPoolingOptions(Element element, BeanDefinitionBuilder builder) { + + addOptionalPropertyValue(builder, "localCoreConnections", element, "core-connections", null); + addOptionalPropertyValue(builder, "localMaxConnections", element, "max-connections", null); + addOptionalPropertyValue(builder, "localMaxSimultaneousRequests", element, "max-simultaneous-requests", null); + addOptionalPropertyValue(builder, "localMinSimultaneousRequests", element, "min-simultaneous-requests", null); + } + + /** + * Parses remote pooling options. + * + * @param element {@link Element} to parse. + * @param builder {@link BeanDefinitionBuilder} used to build a {@link PoolingOptions} {@link BeanDefinition}. + */ + void parseRemotePoolingOptions(Element element, BeanDefinitionBuilder builder) { + + addOptionalPropertyValue(builder, "remoteCoreConnections", element, "core-connections", null); + addOptionalPropertyValue(builder, "remoteMaxConnections", element, "max-connections", null); + addOptionalPropertyValue(builder, "remoteMaxSimultaneousRequests", element, "max-simultaneous-requests", null); + addOptionalPropertyValue(builder, "remoteMinSimultaneousRequests", element, "min-simultaneous-requests", null); + } + + /** + * Parse CQL script {@link Element}s. + * + * @param element {@link Element} to parse. + * @return return the contents of the {@link Element}, which should contain the CQL script. + */ + String parseScript(Element element) { return element.getTextContent(); } - /** - * Returns a {@link BeanDefinition} for a {@link PoolingOptions} object. - * - * @param element The Element being parsed - * @param builder The {@link BeanDefinition} to use for building if one already exists - * @param hostDistance The scope of the PoolingOptions to apply - * @return The {@link BeanDefinitionBuilder} - */ - protected BeanDefinitionBuilder parsePoolingOptions(Element element, BeanDefinitionBuilder builder, - HostDistance hostDistance) { - - if (builder == null) { - builder = BeanDefinitionBuilder.genericBeanDefinition(PoolingOptionsFactoryBean.class); - } - - if (hostDistance.equals(HostDistance.LOCAL)) { - addOptionalPropertyValue(builder, "localMinSimultaneousRequests", element, "min-simultaneous-requests", null); - addOptionalPropertyValue(builder, "localMaxSimultaneousRequests", element, "max-simultaneous-requests", null); - addOptionalPropertyValue(builder, "localCoreConnections", element, "core-connections", null); - addOptionalPropertyValue(builder, "localMaxConnections", element, "max-connections", null); - } - if (hostDistance.equals(HostDistance.REMOTE)) { - addOptionalPropertyValue(builder, "remoteMinSimultaneousRequests", element, "min-simultaneous-requests", null); - addOptionalPropertyValue(builder, "remoteMaxSimultaneousRequests", element, "max-simultaneous-requests", null); - addOptionalPropertyValue(builder, "remoteCoreConnections", element, "core-connections", null); - addOptionalPropertyValue(builder, "remoteMaxConnections", element, "max-connections", null); - } - - return builder; - } - /** * Returns a {@link BeanDefinition} for a {@link SocketOptions} object. - * - * @param element The Element being parsed - * @param context The ParserContext - * @return The {@link BeanDefinition} + * + * @param element {@link Element} to parse. + * @param parserContext XML parser context and state. + * @return {@link BeanDefinition} for {@link SocketOptionsFactoryBean}. */ - protected BeanDefinition getSocketOptionsBeanDefinition(Element element, ParserContext context) { + BeanDefinition newSocketOptionsBeanDefinition(Element element, ParserContext parserContext) { BeanDefinitionBuilder builder = BeanDefinitionBuilder.genericBeanDefinition(SocketOptionsFactoryBean.class); - addOptionalPropertyValue(builder, "connectTimeoutMillis", element, "connect-timeout-mls", null); - addOptionalPropertyValue(builder, "keepAlive", element, "keep-alive", null); - addOptionalPropertyValue(builder, "readTimeoutMillis", element, "read-timeout-mls", null); - addOptionalPropertyValue(builder, "reuseAddress", element, "reuse-address", null); - addOptionalPropertyValue(builder, "soLinger", element, "so-linger", null); - addOptionalPropertyValue(builder, "tcpNoDelay", element, "tcp-no-delay", null); - addOptionalPropertyValue(builder, "receiveBufferSize", element, "receive-buffer-size", null); - addOptionalPropertyValue(builder, "sendBufferSize", element, "send-buffer-size", null); + addOptionalPropertyValue(builder, "connectTimeoutMillis", element, "connect-timeout-millis"); + addOptionalPropertyValue(builder, "keepAlive", element, "keep-alive"); + addOptionalPropertyValue(builder, "readTimeoutMillis", element, "read-timeout-millis"); + addOptionalPropertyValue(builder, "receiveBufferSize", element, "receive-buffer-size"); + addOptionalPropertyValue(builder, "reuseAddress", element, "reuse-address"); + addOptionalPropertyValue(builder, "sendBufferSize", element, "send-buffer-size"); + addOptionalPropertyValue(builder, "soLinger", element, "so-linger"); + addOptionalPropertyValue(builder, "tcpNoDelay", element, "tcp-no-delay"); - return getSourceBeanDefinition(builder, context, element); - } - - /** - * Returns a {@link BeanDefinition} for a {@link KeyspaceActionSpecification} object. - * - * @param element The Element being parsed - * @param context The Parser Context - * @return The {@link BeanDefinition} or {@literal null} if action is not given. - */ - private BeanDefinition getKeyspaceSpecificationBeanDefinition(Element element, ParserContext context) { - - BeanDefinitionBuilder keyspaceBuilder = BeanDefinitionBuilder - .genericBeanDefinition(KeyspaceActionSpecificationFactoryBean.class); - - // add required replication defaults - addRequiredPropertyValue(keyspaceBuilder, "replicationStrategy", - KeyspaceAttributes.DEFAULT_REPLICATION_STRATEGY.name()); - addRequiredPropertyValue(keyspaceBuilder, "replicationFactor", "" + KeyspaceAttributes.DEFAULT_REPLICATION_FACTOR); - - // now start parsing - addRequiredPropertyValue(keyspaceBuilder, "name", element, "name"); - addRequiredPropertyValue(keyspaceBuilder, "action", element, "action"); - addOptionalPropertyValue(keyspaceBuilder, "durableWrites", element, "durable-writes", "false"); - - Element replicationElement = DomUtils.getChildElementByTagName(element, "replication"); - parseReplication(replicationElement, keyspaceBuilder); - - return getSourceBeanDefinition(keyspaceBuilder, context, element); + return getSourceBeanDefinition(builder, parserContext, element); } } diff --git a/spring-cql/src/main/java/org/springframework/cassandra/config/xml/ParsingUtils.java b/spring-cql/src/main/java/org/springframework/cassandra/config/xml/ParsingUtils.java index 97adf3dc8..d6254b82b 100644 --- a/spring-cql/src/main/java/org/springframework/cassandra/config/xml/ParsingUtils.java +++ b/spring-cql/src/main/java/org/springframework/cassandra/config/xml/ParsingUtils.java @@ -1,12 +1,12 @@ /* * Copyright 2013-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. @@ -24,182 +24,69 @@ import org.springframework.util.StringUtils; import org.w3c.dom.Attr; import org.w3c.dom.Element; -public class ParsingUtils { +/** + * Utility class for parsing Cassandra XML namespace configuration meta-data. + * + * @author John Blum + * @see org.springframework.beans.factory.config.BeanDefinition + * @see org.springframework.beans.factory.support.BeanDefinitionBuilder + * @see org.springframework.beans.factory.xml.ParserContext + * @see org.w3c.dom.Attr + * @see org.w3c.dom.Element + */ +public abstract class ParsingUtils { /** - * Convenience method that ultimately delegates to - * {@link #addProperty(BeanDefinitionBuilder, String, Element, String, String, boolean, boolean)}. + * Convenience method delegating to + * {@link #addProperty(BeanDefinitionBuilder, String, String, String, boolean, boolean)}. */ - public static void addOptionalPropertyValue(BeanDefinitionBuilder builder, String propertyName, Element element, - String attrName, String defaultValue) { + public static void addOptionalPropertyReference(BeanDefinitionBuilder builder, String propertyName, Attr attribute) { - addProperty(builder, propertyName, element.getAttribute(attrName), defaultValue, false, false); + addProperty(builder, propertyName, attribute.getValue(), null, false, true); } /** - * Convenience method that ultimately delegates to - * {@link #addProperty(BeanDefinitionBuilder, String, Element, String, String, boolean, boolean)}. + * Convenience method delegating to + * {@link #addProperty(BeanDefinitionBuilder, String, String, String, boolean, boolean)}. + */ + public static void addOptionalPropertyReference(BeanDefinitionBuilder builder, String propertyName, Attr attribute, + String defaultValue) { + + addProperty(builder, propertyName, attribute.getValue(), defaultValue, false, true); + } + + /** + * Convenience method delegating to + * {@link #addProperty(BeanDefinitionBuilder, String, String, String, boolean, boolean)}. */ public static void addOptionalPropertyReference(BeanDefinitionBuilder builder, String propertyName, Element element, - String attrName, String defaultValue) { + String attributeName) { - addProperty(builder, propertyName, element.getAttribute(attrName), defaultValue, false, true); + addProperty(builder, propertyName, element.getAttribute(attributeName), null, false, true); } /** - * Convenience method that ultimately delegates to - * {@link #addProperty(BeanDefinitionBuilder, String, Element, String, String, boolean, boolean)}. + * Convenience method delegating to + * {@link #addProperty(BeanDefinitionBuilder, String, String, String, boolean, boolean)}. */ - public static void addRequiredPropertyValue(BeanDefinitionBuilder builder, String propertyName, Element element, - String attrName) { + public static void addOptionalPropertyReference(BeanDefinitionBuilder builder, String propertyName, Element element, + String attributeName, String defaultValue) { - addProperty(builder, propertyName, element.getAttribute(attrName), null, true, false); + addProperty(builder, propertyName, element.getAttribute(attributeName), defaultValue, false, true); } /** - * Convenience method that ultimately delegates to - * {@link #addProperty(BeanDefinitionBuilder, String, Element, String, String, boolean, boolean)}. + * Convenience method delegating to + * {@link ParsingUtils#addProperty(BeanDefinitionBuilder, String, String, String, boolean, boolean)}. */ - public static void addRequiredPropertyReference(BeanDefinitionBuilder builder, String propertyName, Element element, - String attrName) { + public static void addOptionalPropertyReference(BeanDefinitionBuilder builder, String propertyName, String value) { - addProperty(builder, propertyName, element.getAttribute(attrName), null, true, true); + addProperty(builder, propertyName, value, null, false, true); } /** - * Convenience method that ultimately delegates to - * {@link #addProperty(BeanDefinitionBuilder, String, Element, String, String, boolean, boolean)}. - */ - public static void addPropertyValue(BeanDefinitionBuilder builder, String propertyName, Element element, - String attrName, String defaultValue, boolean required) { - - addProperty(builder, propertyName, element.getAttribute(attrName), defaultValue, required, false); - } - - /** - * Convenience method that ultimately delegates to - * {@link #addProperty(BeanDefinitionBuilder, String, Element, String, String, boolean, boolean)}. - */ - public static void addPropertyReference(BeanDefinitionBuilder builder, String propertyName, Element element, - String attrName, String defaultValue, boolean required) { - - addProperty(builder, propertyName, element.getAttribute(attrName), defaultValue, required, true); - } - - /** - * Convenience method that ultimately delegates to - * {@link #addProperty(BeanDefinitionBuilder, String, Element, String, String, boolean, boolean)}. - */ - public static void addProperty(BeanDefinitionBuilder builder, String propertyName, Element element, String attrName, - String defaultValue, boolean required, boolean reference) { - - Assert.notNull(element, "Element must not be null!"); - Assert.hasText(attrName, "Attribute name must not be null!"); - - addProperty(builder, propertyName, element.getAttribute(attrName), defaultValue, required, reference); - } - - /** - * Convenience method that ultimately delegates to - * {@link #addProperty(BeanDefinitionBuilder, String, Element, String, String, boolean, boolean)}. - */ - public static void addOptionalPropertyValue(BeanDefinitionBuilder builder, String propertyName, Attr attr, - String defaultValue) { - - addProperty(builder, propertyName, attr, defaultValue, false, false); - } - - /** - * Convenience method that ultimately delegates to - * {@link #addProperty(BeanDefinitionBuilder, String, Element, String, String, boolean, boolean)}. - */ - public static void addOptionalPropertyReference(BeanDefinitionBuilder builder, String propertyName, Attr attr, - String defaultValue) { - - addProperty(builder, propertyName, attr, defaultValue, false, true); - } - - /** - * Convenience method that ultimately delegates to - * {@link #addProperty(BeanDefinitionBuilder, String, Element, String, String, boolean, boolean)}. - */ - public static void addRequiredPropertyValue(BeanDefinitionBuilder builder, String propertyName, Attr attr) { - - addProperty(builder, propertyName, attr, null, true, false); - } - - /** - * Convenience method that ultimately delegates to - * {@link #addProperty(BeanDefinitionBuilder, String, Element, String, String, boolean, boolean)}. - */ - public static void addRequiredPropertyReference(BeanDefinitionBuilder builder, String propertyName, Attr attr) { - - addProperty(builder, propertyName, attr, null, true, true); - } - - /** - * Convenience method that ultimately delegates to - * {@link #addProperty(BeanDefinitionBuilder, String, Element, String, String, boolean, boolean)}. - */ - public static void addPropertyValue(BeanDefinitionBuilder builder, String propertyName, Attr attr, - String defaultValue, boolean required) { - - addProperty(builder, propertyName, attr, defaultValue, required, false); - } - - /** - * Convenience method that ultimately delegates to - * {@link #addProperty(BeanDefinitionBuilder, String, Element, String, String, boolean, boolean)}. - */ - public static void addPropertyReference(BeanDefinitionBuilder builder, String propertyName, Attr attr, - String defaultValue, boolean required) { - - addProperty(builder, propertyName, attr, defaultValue, required, true); - } - - /** - * Convenience method that ultimately delegates to - * {@link #addProperty(BeanDefinitionBuilder, String, Element, String, String, boolean, boolean)}. - */ - public static void addProperty(BeanDefinitionBuilder builder, String propertyName, Attr attr, String defaultValue, - boolean required, boolean reference) { - - Assert.notNull(attr, "Attr must not be null!"); - - addProperty(builder, propertyName, attr.getValue(), defaultValue, required, reference); - } - - /** - * Convenience method that ultimately delegates to - * {@link #addProperty(BeanDefinitionBuilder, String, Element, String, String, boolean, boolean)}. - */ - public static void addRequiredPropertyValue(BeanDefinitionBuilder builder, String propertyName, String value) { - - addProperty(builder, propertyName, value, null, true, false); - } - - /** - * Convenience method that ultimately delegates to - * {@link #addProperty(BeanDefinitionBuilder, String, Element, String, String, boolean, boolean)}. - */ - public static void addRequiredPropertyReference(BeanDefinitionBuilder builder, String propertyName, String value) { - - addProperty(builder, propertyName, value, null, true, true); - } - - /** - * Convenience method that ultimately delegates to - * {@link #addProperty(BeanDefinitionBuilder, String, Element, String, String, boolean, boolean)}. - */ - public static void addOptionalPropertyValue(BeanDefinitionBuilder builder, String propertyName, String value, - String defaultValue) { - - addProperty(builder, propertyName, value, defaultValue, false, false); - } - - /** - * Convenience method that ultimately delegates to - * {@link #addProperty(BeanDefinitionBuilder, String, Element, String, String, boolean, boolean)}. + * Convenience method delegating to + * {@link ParsingUtils#addProperty(BeanDefinitionBuilder, String, String, String, boolean, boolean)}. */ public static void addOptionalPropertyReference(BeanDefinitionBuilder builder, String propertyName, String value, String defaultValue) { @@ -208,102 +95,206 @@ public class ParsingUtils { } /** - * Convenience method that ultimately delegates to - * {@link #addProperty(BeanDefinitionBuilder, String, Element, String, String, boolean, boolean)}. + * Convenience method delegating to + * {@link #addProperty(BeanDefinitionBuilder, String, String, String, boolean, boolean)}. */ - public static void addPropertyValue(BeanDefinitionBuilder builder, String propertyName, String value, - String defaultValue, boolean required) { + public static void addOptionalPropertyValue(BeanDefinitionBuilder builder, String propertyName, Attr attribute) { - addProperty(builder, propertyName, value, defaultValue, required, false); + addProperty(builder, propertyName, attribute.getValue(), null, false, false); } /** - * Convenience method that ultimately delegates to - * {@link #addProperty(BeanDefinitionBuilder, String, Element, String, String, boolean, boolean)}. + * Convenience method delegating to + * {@link #addProperty(BeanDefinitionBuilder, String, String, String, boolean, boolean)}. */ - public static void addPropertyReference(BeanDefinitionBuilder builder, String propertyName, String value, - String defaultValue, boolean required) { + public static void addOptionalPropertyValue(BeanDefinitionBuilder builder, String propertyName, Attr attribute, + String defaultValue) { - addProperty(builder, propertyName, value, defaultValue, required, true); + addProperty(builder, propertyName, attribute.getValue(), defaultValue, false, false); } /** - * Adds the named property as a value or reference to the given {@link BeanDefinitionBuilder}, with an optional - * default value. + * Convenience method delegating to + * {@link #addProperty(BeanDefinitionBuilder, String, String, String, boolean, boolean)}. + */ + public static void addOptionalPropertyValue(BeanDefinitionBuilder builder, String propertyName, Element element, + String attributeName) { + + addProperty(builder, propertyName, element.getAttribute(attributeName), null, false, false); + } + + /** + * Convenience method delegating to + * {@link #addProperty(BeanDefinitionBuilder, String, String, String, boolean, boolean)}. + */ + public static void addOptionalPropertyValue(BeanDefinitionBuilder builder, String propertyName, Element element, + String attributeName, String defaultValue) { + + addProperty(builder, propertyName, element.getAttribute(attributeName), defaultValue, false, false); + } + + /** + * Convenience method delegating to + * {@link #addProperty(BeanDefinitionBuilder, String, String, String, boolean, boolean)}. + */ + public static void addOptionalPropertyValue(BeanDefinitionBuilder builder, String propertyName, String value) { + + addProperty(builder, propertyName, value, null, false, false); + } + + /** + * Convenience method delegating to + * {@link #addProperty(BeanDefinitionBuilder, String, String, String, boolean, boolean)}. + */ + public static void addOptionalPropertyValue(BeanDefinitionBuilder builder, String propertyName, String value, + String defaultValue) { + + addProperty(builder, propertyName, value, defaultValue, false, false); + } + + /** + * Convenience method delegating to + * {@link #addProperty(BeanDefinitionBuilder, String, String, String, boolean, boolean)}. + */ + public static void addRequiredPropertyReference(BeanDefinitionBuilder builder, String propertyName, Attr attribute) { + + addProperty(builder, propertyName, attribute.getValue(), null, true, true); + } + + /** + * Convenience method delegating to + * {@link #addProperty(BeanDefinitionBuilder, String, String, String, boolean, boolean)}. + */ + public static void addRequiredPropertyReference(BeanDefinitionBuilder builder, String propertyName, Element element, + String attributeName) { + + addProperty(builder, propertyName, element.getAttribute(attributeName), null, true, true); + } + + /** + * Convenience method delegating to + * {@link #addProperty(BeanDefinitionBuilder, String, String, String, boolean, boolean)}. + */ + public static void addRequiredPropertyReference(BeanDefinitionBuilder builder, String propertyName, String value) { + + addProperty(builder, propertyName, value, null, true, true); + } + + /** + * Convenience method delegating to + * {@link #addProperty(BeanDefinitionBuilder, String, String, String, boolean, boolean)}. + */ + public static void addRequiredPropertyValue(BeanDefinitionBuilder builder, String propertyName, Attr attribute) { + + addProperty(builder, propertyName, attribute.getValue(), null, true, false); + } + + /** + * Convenience method delegating to + * {@link #addProperty(BeanDefinitionBuilder, String, String, String, boolean, boolean)}. + */ + public static void addRequiredPropertyValue(BeanDefinitionBuilder builder, String propertyName, Element element, + String attributeName) { + + addProperty(builder, propertyName, element.getAttribute(attributeName), null, true, false); + } + + /** + * Convenience method delegating to + * {@link #addProperty(BeanDefinitionBuilder, String, String, String, boolean, boolean)}. + */ + public static void addRequiredPropertyValue(BeanDefinitionBuilder builder, String propertyName, String value) { + + addProperty(builder, propertyName, value, null, true, false); + } + + /** + * Adds the named property and value, or reference to the given {@link BeanDefinitionBuilder} with an optional + * default value if the value has not been specified. *

- * Note: If required is false, value is null or empty, and - * defaultValue is null or empty, then no property is added and this method silently returns. - * - * @param builder The {@link BeanDefinitionBuilder}; must not be null. - * @param propertyName The name of the property being added; must not be null or empty. - * @param value The value of the property being added; may be null. - * @param defaultValue The default value of the property being set. - * @param required If true, then the value parameter must not be null or empty. If - * false, the value parameter may be null, in which case the - * defaultValue is used. If required is false, value is - * null or empty, and defaultValue is null or empty, then no property is added and this method - * silently returns. - * @param reference If true, this method will add the property as a reference, else as a value. + * If required is false, value is null or empty, + * and defaultValue is null or empty, then no property is added to the bean definition + * and this method silently returns. + * + * @param builder {@link BeanDefinitionBuilder} used to build the {@link BeanDefinition}. + * @param propertyName name of the property to add. + * @param value value for the property being added. + * @param defaultValue default value for the property if value is null or empty. + * @param required If true, then value must not be null or empty. If false, + * then value may be null and the defaultValue will be used. If required is + * false, value is null or empty, and defaultValue is null or empty, + * then no property is added to the bean definition and this method silently returns. + * @param reference If true, then the valuevalue for the named property is + * considered a reference to another bean in the Spring context. + * @return the given {@link BeanDefinitionBuilder}. + * @throws IllegalArgumentException if either the {@link BeanDefinitionBuilder} is null + * or the propertyName has not been specified. * @see BeanDefinitionBuilder#addPropertyReference(String, String) * @see BeanDefinitionBuilder#addPropertyValue(String, Object) */ - public static void addProperty(BeanDefinitionBuilder builder, String propertyName, String value, String defaultValue, - boolean required, boolean reference) { + public static BeanDefinitionBuilder addProperty(BeanDefinitionBuilder builder, String propertyName, + String value, String defaultValue, boolean required, boolean reference) { - Assert.notNull(builder, "BeanDefinitionBuilder must not be null!"); - Assert.hasText(propertyName, "Property name must not be null!"); + Assert.notNull(builder, "BeanDefinitionBuilder must not be null"); + Assert.hasText(propertyName, "Property name must not be null"); if (!StringUtils.hasText(value)) { if (required) { - throw new IllegalStateException(String.format("value required for property %s [%s] on class [%s]", - reference ? "reference" : "", propertyName, builder.getBeanDefinition().getClass().getName())); + throw new IllegalArgumentException(String.format("value required for property %1$s[%2$s] on class [%3$s]", + reference ? "reference " : "", propertyName, builder.getRawBeanDefinition().getBeanClassName())); } - // else optional; use default - if (defaultValue != null) { + else { value = defaultValue; - } else { // no default value given; quietly ignore & return - return; } } - if (reference) { - builder.addPropertyReference(propertyName, value); - } else { - builder.addPropertyValue(propertyName, value); + if (StringUtils.hasText(value)) { + if (reference) { + builder.addPropertyReference(propertyName, value); + } else { + builder.addPropertyValue(propertyName, value); + } } + + return builder; } /** - * Returns the {@link BeanDefinition} built by the given {@link BeanDefinitionBuilder} enriched with source - * information derived from the given {@link Element}. - * - * @param builder must not be {@literal null}. - * @param context must not be {@literal null}. - * @param element must not be {@literal null}. - * @return + * Returns a {@link BeanDefinition} built from the given {@link BeanDefinitionBuilder} enriched with + * source meta-data derived from the given {@link Element}. + * + * @param builder {@link BeanDefinitionBuilder} used to build the {@link BeanDefinition}. + * @param parserContext {@link ParserContext} used to track state during the parsing operation. + * @param element DOM {@link Element} defining the meta-data that is the source of the {@link BeanDefinition}s + * configuration. + * @return the {@link BeanDefinition} built by the given {@link BeanDefinitionBuilder}. + * @throws IllegalArgumentException if the {@link BeanDefinitionBuilder} or {@link ParserContext} are null. */ - public static AbstractBeanDefinition getSourceBeanDefinition(BeanDefinitionBuilder builder, ParserContext context, - Element element) { + public static AbstractBeanDefinition getSourceBeanDefinition(BeanDefinitionBuilder builder, + ParserContext parserContext, Element element) { - Assert.notNull(element, "Element must not be null!"); - Assert.notNull(context, "ParserContext must not be null!"); + Assert.notNull(parserContext, "ParserContext must not be null"); - return getSourceBeanDefinition(builder, context.extractSource(element)); + return getSourceBeanDefinition(builder, parserContext.extractSource(element)); } /** - * Returns the {@link AbstractBeanDefinition} built by the given builder with the given extracted source applied. - * - * @param builder must not be {@literal null}. - * @param source - * @return + * Returns a {@link AbstractBeanDefinition} built from the given {@link BeanDefinitionBuilder} with the given + * extracted source applied. + * + * @param builder {@link BeanDefinitionBuilder} used to build the {@link BeanDefinition}. + * @param source source meta-data used by the builder to construct the {@link BeanDefinition}. + * @return a raw {@link BeanDefinition} built by the given {@link BeanDefinitionBuilder}. + * @throws IllegalArgumentException if {@link BeanDefinitionBuilder} is null. */ public static AbstractBeanDefinition getSourceBeanDefinition(BeanDefinitionBuilder builder, Object source) { - Assert.notNull(builder, "Builder must not be null!"); + Assert.notNull(builder, "BeanDefinitionBuilder must not be null"); - AbstractBeanDefinition definition = builder.getRawBeanDefinition(); - definition.setSource(source); - return definition; + AbstractBeanDefinition beanDefinition = builder.getBeanDefinition(); + + beanDefinition.setSource(source); + + return beanDefinition; } } diff --git a/spring-cql/src/main/resources/org/springframework/cassandra/config/spring-cql-1.0.xsd b/spring-cql/src/main/resources/org/springframework/cassandra/config/spring-cql-1.0.xsd index c7cece5c2..c51c2f7b9 100644 --- a/spring-cql/src/main/resources/org/springframework/cassandra/config/spring-cql-1.0.xsd +++ b/spring-cql/src/main/resources/org/springframework/cassandra/config/spring-cql-1.0.xsd @@ -1,47 +1,30 @@ - + xmlns:tool="http://www.springframework.org/schema/tool" + xmlns:xsd="http://www.w3.org/2001/XMLSchema" + xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" + targetNamespace="http://www.springframework.org/schema/cql" + elementFormDefault="qualified" attributeFormDefault="unqualified"> - + - - - + + - - - - - - - - - - - - + + @@ -57,42 +40,48 @@ Defines a Cassandra cluster. + + + + + + + + + + + - + - + - + - + - + - + - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + - + - - - - - - - - - - - - - - - - - - - - - - - - - - + + + - - + + - - + @@ -325,21 +347,141 @@ Custom Latency Tracker for the Cassandra Cluster. + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + - - - - - - - - - - + + + + + + + + + + @@ -378,6 +534,13 @@ Sets read timeout for client socket in milliseconds. ]]> + + + + + + + + + + - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - @@ -525,8 +578,7 @@ Provides the ability to configure the keyspace's replication settings. ]]> - + - + - + - - - - - - - - - - - - - - - - \ No newline at end of file + + diff --git a/spring-cql/src/test/java/org/springframework/cassandra/config/PoolingOptionsFactoryBeanUnitTest.java b/spring-cql/src/test/java/org/springframework/cassandra/config/PoolingOptionsFactoryBeanUnitTest.java deleted file mode 100644 index 57818a123..000000000 --- a/spring-cql/src/test/java/org/springframework/cassandra/config/PoolingOptionsFactoryBeanUnitTest.java +++ /dev/null @@ -1,70 +0,0 @@ -/* - * Copyright 2016 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.config; - -import org.junit.Assert; -import org.junit.Test; - -/** - * Unit tests for {@link PoolingOptionsFactoryBean} - * - * @author Sumit Kumar - * @author David Webb - */ -public class PoolingOptionsFactoryBeanUnitTest { - - private static final int REMOTE_MIN_SIMULTANEOUS_REQUESTS = 111; - private static final int REMOTE_MAX_SIMULTANEOUS_REQUESTS = 127; - private static final int REMOTE_CORE_CONNECTIONS = 110; - private static final int REMOTE_MAX_CONNECTIONS = 210; - private static final int LOCAL_MIN_SIMULTANEOUS_REQUESTS = 97; - private static final int LOCAL_MAX_SIMULTANEOUS_REQUESTS = 99; - private static final int LOCAL_CORE_CONNECTIONS = 100; - private static final int LOCAL_MAX_CONNECTIONS = 200; - - /** - * The max values should be set before setting core values. Otherwise the core values will be compared with the - * default max values which is 8. Same for other min-max properties pairs. This test checks the same. - * - * @throws Exception Any unhandled scenarios will result in a test failure. - * @see DATACASS-176 - */ - @Test - public void testAfterPropertiesSet() throws Exception { - - PoolingOptionsFactoryBean factoryBean = new PoolingOptionsFactoryBean(); - factoryBean.setLocalMaxConnections(LOCAL_MAX_CONNECTIONS); - factoryBean.setLocalCoreConnections(LOCAL_CORE_CONNECTIONS); - factoryBean.setLocalMaxSimultaneousRequests(LOCAL_MAX_SIMULTANEOUS_REQUESTS); - factoryBean.setLocalMinSimultaneousRequests(LOCAL_MIN_SIMULTANEOUS_REQUESTS); - factoryBean.setRemoteMaxConnections(REMOTE_MAX_CONNECTIONS); - factoryBean.setRemoteCoreConnections(REMOTE_CORE_CONNECTIONS); - factoryBean.setRemoteMaxSimultaneousRequests(REMOTE_MAX_SIMULTANEOUS_REQUESTS); - factoryBean.setRemoteMinSimultaneousRequests(REMOTE_MIN_SIMULTANEOUS_REQUESTS); - factoryBean.afterPropertiesSet(); - - Assert.assertEquals(factoryBean.getLocalMaxConnections().intValue(), LOCAL_MAX_CONNECTIONS); - Assert.assertEquals(factoryBean.getLocalCoreConnections().intValue(), LOCAL_CORE_CONNECTIONS); - Assert.assertEquals(factoryBean.getLocalMaxSimultaneousRequests().intValue(), LOCAL_MAX_SIMULTANEOUS_REQUESTS); - Assert.assertEquals(factoryBean.getLocalMinSimultaneousRequests().intValue(), LOCAL_MIN_SIMULTANEOUS_REQUESTS); - Assert.assertEquals(factoryBean.getRemoteMaxConnections().intValue(), REMOTE_MAX_CONNECTIONS); - Assert.assertEquals(factoryBean.getRemoteCoreConnections().intValue(), REMOTE_CORE_CONNECTIONS); - Assert.assertEquals(factoryBean.getRemoteMaxSimultaneousRequests().intValue(), REMOTE_MAX_SIMULTANEOUS_REQUESTS); - Assert.assertEquals(factoryBean.getRemoteMinSimultaneousRequests().intValue(), REMOTE_MIN_SIMULTANEOUS_REQUESTS); - - } - -} diff --git a/spring-cql/src/test/java/org/springframework/cassandra/config/PoolingOptionsFactoryBeanUnitTests.java b/spring-cql/src/test/java/org/springframework/cassandra/config/PoolingOptionsFactoryBeanUnitTests.java new file mode 100644 index 000000000..ec8b374e0 --- /dev/null +++ b/spring-cql/src/test/java/org/springframework/cassandra/config/PoolingOptionsFactoryBeanUnitTests.java @@ -0,0 +1,347 @@ +/* + * Copyright 2016 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.config; + +import static org.hamcrest.Matchers.*; +import static org.junit.Assert.*; +import static org.mockito.Matchers.eq; +import static org.mockito.Mockito.any; +import static org.mockito.Mockito.anyInt; +import static org.mockito.Mockito.*; +import static org.mockito.Mockito.same; + +import java.util.concurrent.Executor; + +import org.junit.Before; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.InOrder; +import org.mockito.Mock; +import org.mockito.Spy; +import org.mockito.invocation.InvocationOnMock; +import org.mockito.runners.MockitoJUnitRunner; +import org.mockito.stubbing.Answer; + +import com.datastax.driver.core.HostDistance; +import com.datastax.driver.core.PoolingOptions; + +/** + * Unit tests for {@link PoolingOptionsFactoryBean}. + * + * @author Sumit Kumar + * @author David Webb + * @author John Blum + * @see org.springframework.cassandra.config.PoolingOptionsFactoryBean + * @see DATACASS-176 + * @see DATACASS-298 + */ +@RunWith(MockitoJUnitRunner.class) +public class PoolingOptionsFactoryBeanUnitTests { + + @Mock + private Executor mockExecutor; + + @Spy + private PoolingOptions poolingOptionsSpy; + + private PoolingOptionsFactoryBean poolingOptionsFactoryBean; + + @Before + public void setup() { + poolingOptionsFactoryBean = new PoolingOptionsFactoryBean(); + } + + @Test + public void getObjectReturnsNullWhenNotInitialized() throws Exception { + assertThat(poolingOptionsFactoryBean.getObject(), is(nullValue(PoolingOptions.class))); + } + + @Test + @SuppressWarnings("unchecked") + public void getObjectTypeReturnsPoolingOptionsClassWhenNotInitialized() { + assertThat((Class) poolingOptionsFactoryBean.getObjectType(), is(equalTo(PoolingOptions.class))); + } + + @Test + public void isSingletonIsTrue() { + assertThat(poolingOptionsFactoryBean.isSingleton(), is(true)); + } + + /** + * @see DATACASS-298 + */ + @Test + public void setAndGetFactoryBeanProperties() { + poolingOptionsFactoryBean.setHeartbeatIntervalSeconds(15); + poolingOptionsFactoryBean.setIdleTimeoutSeconds(120); + poolingOptionsFactoryBean.setInitializationExecutor(mockExecutor); + poolingOptionsFactoryBean.setLocalCoreConnections(50); + poolingOptionsFactoryBean.setLocalMaxConnections(1000); + poolingOptionsFactoryBean.setLocalMaxSimultaneousRequests(200); + poolingOptionsFactoryBean.setLocalMinSimultaneousRequests(100); + poolingOptionsFactoryBean.setPoolTimeoutMilliseconds(300); + poolingOptionsFactoryBean.setRemoteCoreConnections(25); + poolingOptionsFactoryBean.setRemoteMaxConnections(250); + poolingOptionsFactoryBean.setRemoteMaxSimultaneousRequests(100); + poolingOptionsFactoryBean.setRemoteMinSimultaneousRequests(50); + + assertThat(poolingOptionsFactoryBean.getHeartbeatIntervalSeconds(), is(equalTo(15))); + assertThat(poolingOptionsFactoryBean.getIdleTimeoutSeconds(), is(equalTo(120))); + assertThat(poolingOptionsFactoryBean.getInitializationExecutor(), is(equalTo(mockExecutor))); + assertThat(poolingOptionsFactoryBean.getLocalCoreConnections(), is(equalTo(50))); + assertThat(poolingOptionsFactoryBean.getLocalMaxConnections(), is(equalTo(1000))); + assertThat(poolingOptionsFactoryBean.getLocalMaxSimultaneousRequests(), is(equalTo(200))); + assertThat(poolingOptionsFactoryBean.getLocalMinSimultaneousRequests(), is(equalTo(100))); + assertThat(poolingOptionsFactoryBean.getPoolTimeoutMilliseconds(), is(equalTo(300))); + assertThat(poolingOptionsFactoryBean.getRemoteCoreConnections(), is(equalTo(25))); + assertThat(poolingOptionsFactoryBean.getRemoteMaxConnections(), is(equalTo(250))); + assertThat(poolingOptionsFactoryBean.getRemoteMaxSimultaneousRequests(), is(equalTo(100))); + assertThat(poolingOptionsFactoryBean.getRemoteMinSimultaneousRequests(), is(equalTo(50))); + } + + /** + * @see DATACASS-298 + */ + @Test + public void afterPropertiesSetInitializesLocalPoolingOptions() throws Exception { + PoolingOptionsFactoryBean poolingOptionsFactoryBean = new PoolingOptionsFactoryBean() { + @Override PoolingOptions newPoolingOptions() { + poolingOptionsSpy.setNewConnectionThreshold(HostDistance.LOCAL, 1); + return poolingOptionsSpy; + } + }; + + poolingOptionsFactoryBean.setHeartbeatIntervalSeconds(60); + poolingOptionsFactoryBean.setIdleTimeoutSeconds(300); + poolingOptionsFactoryBean.setInitializationExecutor(mockExecutor); + poolingOptionsFactoryBean.setLocalCoreConnections(10); + poolingOptionsFactoryBean.setLocalMaxConnections(100); + poolingOptionsFactoryBean.setLocalMaxSimultaneousRequests(50); + poolingOptionsFactoryBean.setLocalMinSimultaneousRequests(5); + poolingOptionsFactoryBean.setPoolTimeoutMilliseconds(180); + + assertThat(poolingOptionsFactoryBean.getObject(), is(nullValue(PoolingOptions.class))); + + poolingOptionsFactoryBean.afterPropertiesSet(); + + assertThat(poolingOptionsFactoryBean.getObject(), is(sameInstance(poolingOptionsSpy))); + assertThat(poolingOptionsFactoryBean.getObjectType(), is(equalTo((Class) poolingOptionsSpy.getClass()))); + + verify(poolingOptionsSpy, times(1)).setHeartbeatIntervalSeconds(eq(60)); + verify(poolingOptionsSpy, times(1)).setIdleTimeoutSeconds(eq(300)); + verify(poolingOptionsSpy, times(1)).setInitializationExecutor(eq(mockExecutor)); + verify(poolingOptionsSpy, times(1)).setPoolTimeoutMillis(eq(180)); + verify(poolingOptionsSpy, times(1)).setCoreConnectionsPerHost(eq(HostDistance.LOCAL), eq(10)); + verify(poolingOptionsSpy, times(1)).setMaxConnectionsPerHost(eq(HostDistance.LOCAL), eq(100)); + verify(poolingOptionsSpy, times(1)).setMaxRequestsPerConnection(eq(HostDistance.LOCAL), eq(50)); + verify(poolingOptionsSpy, times(1)).setNewConnectionThreshold(eq(HostDistance.LOCAL), eq(5)); + verify(poolingOptionsSpy, never()).setCoreConnectionsPerHost(eq(HostDistance.REMOTE), anyInt()); + verify(poolingOptionsSpy, never()).setMaxConnectionsPerHost(eq(HostDistance.REMOTE), anyInt()); + verify(poolingOptionsSpy, never()).setMaxRequestsPerConnection(eq(HostDistance.REMOTE), anyInt()); + verify(poolingOptionsSpy, never()).setNewConnectionThreshold(eq(HostDistance.REMOTE), anyInt()); + } + + /** + * @see DATACASS-298 + */ + @Test + public void afterPropertiesSetInitializesRemotePoolingOptions() throws Exception { + PoolingOptionsFactoryBean poolingOptionsFactoryBean = new PoolingOptionsFactoryBean() { + @Override PoolingOptions newPoolingOptions() { + poolingOptionsSpy.setNewConnectionThreshold(HostDistance.REMOTE, 10); + return poolingOptionsSpy; + } + }; + + poolingOptionsFactoryBean.setHeartbeatIntervalSeconds(30); + poolingOptionsFactoryBean.setIdleTimeoutSeconds(120); + poolingOptionsFactoryBean.setInitializationExecutor(mockExecutor); + poolingOptionsFactoryBean.setPoolTimeoutMilliseconds(120); + poolingOptionsFactoryBean.setRemoteCoreConnections(5); + poolingOptionsFactoryBean.setRemoteMaxConnections(50); + poolingOptionsFactoryBean.setRemoteMaxSimultaneousRequests(20); + poolingOptionsFactoryBean.setRemoteMinSimultaneousRequests(5); + + assertThat(poolingOptionsFactoryBean.getObject(), is(nullValue(PoolingOptions.class))); + + poolingOptionsFactoryBean.afterPropertiesSet(); + + assertThat(poolingOptionsFactoryBean.getObject(), is(sameInstance(poolingOptionsSpy))); + assertThat(poolingOptionsFactoryBean.getObjectType(), is(equalTo((Class) poolingOptionsSpy.getClass()))); + + verify(poolingOptionsSpy, times(1)).setHeartbeatIntervalSeconds(eq(30)); + verify(poolingOptionsSpy, times(1)).setIdleTimeoutSeconds(eq(120)); + verify(poolingOptionsSpy, times(1)).setInitializationExecutor(eq(mockExecutor)); + verify(poolingOptionsSpy, times(1)).setPoolTimeoutMillis(eq(120)); + verify(poolingOptionsSpy, times(1)).setCoreConnectionsPerHost(eq(HostDistance.REMOTE), eq(5)); + verify(poolingOptionsSpy, times(1)).setMaxConnectionsPerHost(eq(HostDistance.REMOTE), eq(50)); + verify(poolingOptionsSpy, times(1)).setMaxRequestsPerConnection(eq(HostDistance.REMOTE), eq(20)); + verify(poolingOptionsSpy, never()).setCoreConnectionsPerHost(eq(HostDistance.LOCAL), anyInt()); + verify(poolingOptionsSpy, never()).setMaxConnectionsPerHost(eq(HostDistance.LOCAL), anyInt()); + verify(poolingOptionsSpy, never()).setMaxRequestsPerConnection(eq(HostDistance.LOCAL), anyInt()); + verify(poolingOptionsSpy, never()).setNewConnectionThreshold(eq(HostDistance.LOCAL), anyInt()); + verify(poolingOptionsSpy, never()).setNewConnectionThreshold(eq(HostDistance.REMOTE), eq(5)); + } + + /** + * This particular test case is technically an integration test since it uses an actual instance of + * a DataStax Java driver class type... {@link PoolingOptions}! + * + * The max values should be set before setting core values. Otherwise the core values will be compared with the + * default max values which is 8. Same for other min-max properties pairs. This test checks the same. + * + * @throws Exception Any unhandled scenarios will result in a test failure. + * @see DATACASS-176 + */ + @Test + public void afterPropertiesSetProperlySetsPoolingOptionsMaxBeforeMinProperties() throws Exception { + poolingOptionsFactoryBean = new PoolingOptionsFactoryBean() { + @Override PoolingOptions newPoolingOptions() { + return spy(super.newPoolingOptions()); + } + }; + + poolingOptionsFactoryBean.setLocalMaxConnections(200); + poolingOptionsFactoryBean.setLocalCoreConnections(100); + poolingOptionsFactoryBean.setLocalMaxSimultaneousRequests(99); + poolingOptionsFactoryBean.setLocalMinSimultaneousRequests(97); + poolingOptionsFactoryBean.setRemoteMaxConnections(210); + poolingOptionsFactoryBean.setRemoteCoreConnections(110); + poolingOptionsFactoryBean.setRemoteMaxSimultaneousRequests(127); + poolingOptionsFactoryBean.setRemoteMinSimultaneousRequests(111); + + assertThat(poolingOptionsFactoryBean.getObject(), is(nullValue(PoolingOptions.class))); + + poolingOptionsFactoryBean.afterPropertiesSet(); + + PoolingOptions poolingOptions = poolingOptionsFactoryBean.getObject(); + + assertThat(poolingOptions, is(notNullValue(PoolingOptions.class))); + assertThat(poolingOptions.getCoreConnectionsPerHost(HostDistance.LOCAL), is(equalTo(100))); + assertThat(poolingOptions.getMaxConnectionsPerHost(HostDistance.LOCAL), is(equalTo(200))); + assertThat(poolingOptions.getMaxRequestsPerConnection(HostDistance.LOCAL), is(equalTo(99))); + assertThat(poolingOptions.getNewConnectionThreshold(HostDistance.LOCAL), is(equalTo(97))); + assertThat(poolingOptions.getCoreConnectionsPerHost(HostDistance.REMOTE), is(equalTo(110))); + assertThat(poolingOptions.getMaxConnectionsPerHost(HostDistance.REMOTE), is(equalTo(210))); + assertThat(poolingOptions.getMaxRequestsPerConnection(HostDistance.REMOTE), is(equalTo(127))); + assertThat(poolingOptions.getNewConnectionThreshold(HostDistance.REMOTE), is(equalTo(111))); + + InOrder inOrder = inOrder(poolingOptions); + + inOrder.verify(poolingOptions, times(1)).setMaxConnectionsPerHost(eq(HostDistance.LOCAL), eq(200)); + inOrder.verify(poolingOptions, times(1)).setCoreConnectionsPerHost(eq(HostDistance.LOCAL), eq(100)); + inOrder.verify(poolingOptions, times(1)).setMaxRequestsPerConnection(eq(HostDistance.LOCAL), eq(99)); + inOrder.verify(poolingOptions, times(1)).setNewConnectionThreshold(eq(HostDistance.LOCAL), eq(97)); + inOrder.verify(poolingOptions, times(1)).setMaxConnectionsPerHost(eq(HostDistance.REMOTE), eq(210)); + inOrder.verify(poolingOptions, times(1)).setCoreConnectionsPerHost(eq(HostDistance.REMOTE), eq(110)); + inOrder.verify(poolingOptions, times(1)).setMaxRequestsPerConnection(eq(HostDistance.REMOTE), eq(127)); + inOrder.verify(poolingOptions, times(1)).setNewConnectionThreshold(eq(HostDistance.REMOTE), eq(111)); + } + + @Test + public void newLocalHostDistancePoolingOptionsReturnsLocalHostDistancePoolingOptionsFactoryBeanSettings() { + poolingOptionsFactoryBean.setLocalCoreConnections(50); + poolingOptionsFactoryBean.setLocalMaxConnections(500); + poolingOptionsFactoryBean.setLocalMaxSimultaneousRequests(1000); + poolingOptionsFactoryBean.setLocalMinSimultaneousRequests(100); + poolingOptionsFactoryBean.setRemoteCoreConnections(20); + poolingOptionsFactoryBean.setRemoteMaxConnections(200); + poolingOptionsFactoryBean.setRemoteMaxSimultaneousRequests(400); + poolingOptionsFactoryBean.setRemoteMinSimultaneousRequests(40); + + PoolingOptionsFactoryBean.HostDistancePoolingOptions poolingOptions = + poolingOptionsFactoryBean.newLocalHostDistancePoolingOptions(); + + assertThat(poolingOptions.getHostDistance(), is(equalTo(HostDistance.LOCAL))); + assertThat(poolingOptions.getCoreConnectionsPerHost(), is(equalTo(50))); + assertThat(poolingOptions.getMaxConnectionsPerHost(), is(equalTo(500))); + assertThat(poolingOptions.getMaxRequestsPerConnection(), is(equalTo(1000))); + assertThat(poolingOptions.getNewConnectionThreshold(), is(equalTo(100))); + } + + @Test + public void newLocalHostDistancePoolingOptionsReturnsRemoteHostDistancePoolingOptionsFactoryBeanSettings() { + poolingOptionsFactoryBean.setLocalCoreConnections(50); + poolingOptionsFactoryBean.setLocalMaxConnections(500); + poolingOptionsFactoryBean.setLocalMaxSimultaneousRequests(1000); + poolingOptionsFactoryBean.setLocalMinSimultaneousRequests(100); + poolingOptionsFactoryBean.setRemoteCoreConnections(20); + poolingOptionsFactoryBean.setRemoteMaxConnections(200); + poolingOptionsFactoryBean.setRemoteMaxSimultaneousRequests(400); + poolingOptionsFactoryBean.setRemoteMinSimultaneousRequests(40); + + PoolingOptionsFactoryBean.HostDistancePoolingOptions poolingOptions = + poolingOptionsFactoryBean.newRemoteHostDistancePoolingOptions(); + + assertThat(poolingOptions.getHostDistance(), is(equalTo(HostDistance.REMOTE))); + assertThat(poolingOptions.getCoreConnectionsPerHost(), is(equalTo(20))); + assertThat(poolingOptions.getMaxConnectionsPerHost(), is(equalTo(200))); + assertThat(poolingOptions.getMaxRequestsPerConnection(), is(equalTo(400))); + assertThat(poolingOptions.getNewConnectionThreshold(), is(equalTo(40))); + } + + @Test + public void configureLocalHostDistancePoolingOptionsCallsConfigureWithExpectedInstance() { + final PoolingOptionsFactoryBean.HostDistancePoolingOptions mockHostDistancePoolingOptions = mock( + PoolingOptionsFactoryBean.HostDistancePoolingOptions.class); + + when(mockHostDistancePoolingOptions.configure(any(PoolingOptions.class))).thenAnswer( + new Answer() { + @Override + public PoolingOptions answer(InvocationOnMock invocationOnMock) throws Throwable { + return invocationOnMock.getArgumentAt(0, PoolingOptions.class); + } + } + ); + + poolingOptionsFactoryBean = new PoolingOptionsFactoryBean() { + @Override protected HostDistancePoolingOptions newLocalHostDistancePoolingOptions() { + return mockHostDistancePoolingOptions; + } + }; + + assertThat(poolingOptionsFactoryBean.configureLocalHostDistancePoolingOptions(poolingOptionsSpy), + is(sameInstance(poolingOptionsSpy))); + + verify(mockHostDistancePoolingOptions, times(1)).configure(same(poolingOptionsSpy)); + } + + @Test + public void configureRemoteHostDistancePoolingOptionsCallsConfigureWithExpectedInstance() { + final PoolingOptionsFactoryBean.HostDistancePoolingOptions mockHostDistancePoolingOptions = mock( + PoolingOptionsFactoryBean.HostDistancePoolingOptions.class); + + when(mockHostDistancePoolingOptions.configure(any(PoolingOptions.class))).thenAnswer( + new Answer() { + @Override + public PoolingOptions answer(InvocationOnMock invocationOnMock) throws Throwable { + return invocationOnMock.getArgumentAt(0, PoolingOptions.class); + } + } + ); + + poolingOptionsFactoryBean = new PoolingOptionsFactoryBean() { + @Override protected HostDistancePoolingOptions newRemoteHostDistancePoolingOptions() { + return mockHostDistancePoolingOptions; + } + }; + + assertThat(poolingOptionsFactoryBean.configureRemoteHostDistancePoolingOptions(poolingOptionsSpy), + is(sameInstance(poolingOptionsSpy))); + + verify(mockHostDistancePoolingOptions, times(1)).configure(same(poolingOptionsSpy)); + } +} diff --git a/spring-cql/src/test/java/org/springframework/cassandra/config/xml/CassandraCqlClusterParserUnitTests.java b/spring-cql/src/test/java/org/springframework/cassandra/config/xml/CassandraCqlClusterParserUnitTests.java new file mode 100644 index 000000000..cd3be124a --- /dev/null +++ b/spring-cql/src/test/java/org/springframework/cassandra/config/xml/CassandraCqlClusterParserUnitTests.java @@ -0,0 +1,403 @@ +/* + * Copyright 2013-2016 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.config.xml; + +import static org.hamcrest.Matchers.*; +import static org.junit.Assert.*; +import static org.mockito.Mockito.*; + +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.Mock; +import org.mockito.runners.MockitoJUnitRunner; +import org.springframework.beans.PropertyValue; +import org.springframework.beans.factory.config.BeanDefinition; +import org.springframework.beans.factory.config.RuntimeBeanReference; +import org.springframework.beans.factory.parsing.PassThroughSourceExtractor; +import org.springframework.beans.factory.support.AbstractBeanDefinition; +import org.springframework.beans.factory.support.BeanDefinitionBuilder; +import org.springframework.beans.factory.xml.BeanDefinitionParserDelegate; +import org.springframework.beans.factory.xml.ParserContext; +import org.springframework.beans.factory.xml.XmlReaderContext; +import org.springframework.cassandra.config.CassandraCqlClusterFactoryBean; +import org.springframework.cassandra.config.PoolingOptionsFactoryBean; +import org.springframework.cassandra.config.SocketOptionsFactoryBean; +import org.w3c.dom.Element; +import org.w3c.dom.NodeList; + +/** + * Test suite of Unit tests testing the contract and functionality of the {@link CassandraCqlClusterParser}. + * + * @author John Blum + * @see org.springframework.cassandra.config.xml.CassandraCqlClusterParser + * @since 1.5.0 + */ +// TODO add more tests! +@RunWith(MockitoJUnitRunner.class) +public class CassandraCqlClusterParserUnitTests { + + @Mock + private Element mockElement; + + private CassandraCqlClusterParser parser = new CassandraCqlClusterParser(); + + @SuppressWarnings("unchecked") + protected T getPropertyValue(BeanDefinition beanDefinition, String propertyName) { + PropertyValue propertyValue = beanDefinition.getPropertyValues().getPropertyValue(propertyName); + + return (T) (propertyValue != null ? propertyValue.getValue() : null); + } + + protected String getPropertyValueAsString(BeanDefinition beanDefinition, String propertyName) { + Object value = getPropertyValue(beanDefinition, propertyName); + + return (value instanceof RuntimeBeanReference ? ((RuntimeBeanReference) value).getBeanName() + : (value != null ? String.valueOf(value) : null)); + } + + protected BeanDefinitionParserDelegate mockBeanDefinitionParserDelegate(XmlReaderContext xmlReaderContext) { + return new BeanDefinitionParserDelegate(xmlReaderContext); + } + + protected NodeList mockNodeList(Element... childElements) { + NodeList mockNodeList = mock(NodeList.class); + + when(mockNodeList.getLength()).thenReturn(childElements.length); + + for (int index = 0; index < childElements.length; index++) { + when(mockNodeList.item(eq(index))).thenReturn(childElements[index]); + } + + return mockNodeList; + } + + protected ParserContext mockParserContext() { + return mockParserContext(null); + } + + protected ParserContext mockParserContext(BeanDefinition beanDefinition) { + XmlReaderContext readerContext = mockXmlReaderContext(); + return new ParserContext(readerContext, mockBeanDefinitionParserDelegate(readerContext), beanDefinition); + } + + protected XmlReaderContext mockXmlReaderContext() { + return new XmlReaderContext(null, null, null, new PassThroughSourceExtractor(), null, null); + } + + @Test + public void resolveIdFromElement() { + when(mockElement.getAttribute(eq(CassandraCqlClusterParser.ID_ATTRIBUTE))).thenReturn("test"); + assertThat(parser.resolveId(mockElement, null, null), is(equalTo("test"))); + verify(mockElement, times(1)).getAttribute(eq(CassandraCqlClusterParser.ID_ATTRIBUTE)); + } + + @Test + public void resolveIdUsingDefault() { + when(mockElement.getAttribute(eq(CassandraCqlClusterParser.ID_ATTRIBUTE))).thenReturn(""); + assertThat(parser.resolveId(mockElement, null, null), is(equalTo(DefaultCqlBeanNames.CLUSTER))); + verify(mockElement, times(1)).getAttribute(eq(CassandraCqlClusterParser.ID_ATTRIBUTE)); + } + + @Test + public void parseInternalCallsDoParseAndConstructsBeanDefinition() { + BeanDefinition mockContainingBeanDefinition = mock(BeanDefinition.class); + + when(mockContainingBeanDefinition.getScope()).thenReturn("Singleton"); + when(mockElement.getAttribute("auth-info-provider-ref")).thenReturn("testAuthInfoProvider"); + when(mockElement.getAttribute("host-state-listener-ref")).thenReturn("testHostStateListener"); + when(mockElement.getAttribute("latency-tracker-ref")).thenReturn("testLatencyTracker"); + when(mockElement.getAttribute("load-balancing-policy-ref")).thenReturn("testLoadBalancingPolicy"); + when(mockElement.getAttribute("reconnection-policy-ref")).thenReturn("testReconnectionPolicy"); + when(mockElement.getAttribute("retry-policy-ref")).thenReturn("testRetryPolicy"); + when(mockElement.getAttribute("ssl-options-ref")).thenReturn("testSslOptions"); + when(mockElement.getAttribute("contact-points")).thenReturn("skullbox"); + when(mockElement.getAttribute("compression")).thenReturn("SNAPPY"); + when(mockElement.getAttribute("jmx-reporting-enabled")).thenReturn("true"); + when(mockElement.getAttribute("metrics-enabled")).thenReturn("true"); + when(mockElement.getAttribute("password")).thenReturn("p@55w0rd"); + when(mockElement.getAttribute("port")).thenReturn("12345"); + when(mockElement.getAttribute("ssl-enabled")).thenReturn("true"); + when(mockElement.getAttribute("username")).thenReturn("jonDoe"); + + CassandraCqlClusterParser parser = new CassandraCqlClusterParser() { + @Override + protected void parseChildElements(Element element, ParserContext parserContext, + BeanDefinitionBuilder builder) { + } + }; + + AbstractBeanDefinition beanDefinition = parser.parseInternal(mockElement, mockParserContext( + mockContainingBeanDefinition)); + + assertThat(beanDefinition, is(notNullValue(BeanDefinition.class))); + assertThat(beanDefinition.getBeanClassName(), is(equalTo(CassandraCqlClusterFactoryBean.class.getName()))); + assertThat(beanDefinition.getDestroyMethodName(), is(equalTo("destroy"))); + assertThat((Element) beanDefinition.getSource(), is(equalTo(mockElement))); + assertThat(beanDefinition.isLazyInit(), is(false)); + assertThat(getPropertyValueAsString(beanDefinition, "authProvider"), is(equalTo("testAuthInfoProvider"))); + assertThat(getPropertyValueAsString(beanDefinition, "hostStateListener"), is(equalTo("testHostStateListener"))); + assertThat(getPropertyValueAsString(beanDefinition, "latencyTracker"), is(equalTo("testLatencyTracker"))); + assertThat(getPropertyValueAsString(beanDefinition, "loadBalancingPolicy"), is(equalTo("testLoadBalancingPolicy"))); + assertThat(getPropertyValueAsString(beanDefinition, "reconnectionPolicy"), is(equalTo("testReconnectionPolicy"))); + assertThat(getPropertyValueAsString(beanDefinition, "retryPolicy"), is(equalTo("testRetryPolicy"))); + assertThat(getPropertyValueAsString(beanDefinition, "sslOptions"), is(equalTo("testSslOptions"))); + assertThat(getPropertyValueAsString(beanDefinition, "contactPoints"), is(equalTo("skullbox"))); + assertThat(getPropertyValueAsString(beanDefinition, "compressionType"), is(equalTo("SNAPPY"))); + assertThat(getPropertyValueAsString(beanDefinition, "jmxReportingEnabled"), is(equalTo("true"))); + assertThat(getPropertyValueAsString(beanDefinition, "metricsEnabled"), is(equalTo("true"))); + assertThat(getPropertyValueAsString(beanDefinition, "password"), is(equalTo("p@55w0rd"))); + assertThat(getPropertyValueAsString(beanDefinition, "port"), is(equalTo("12345"))); + assertThat(getPropertyValueAsString(beanDefinition, "sslEnabled"), is(equalTo("true"))); + assertThat(getPropertyValueAsString(beanDefinition, "username"), is(equalTo("jonDoe"))); + + verify(mockContainingBeanDefinition, times(1)).getScope(); + verify(mockElement, times(1)).getAttribute(eq("auth-info-provider-ref")); + verify(mockElement, times(1)).getAttribute(eq("host-state-listener-ref")); + verify(mockElement, times(1)).getAttribute(eq("latency-tracker-ref")); + verify(mockElement, times(1)).getAttribute(eq("load-balancing-policy-ref")); + verify(mockElement, times(1)).getAttribute(eq("reconnection-policy-ref")); + verify(mockElement, times(1)).getAttribute(eq("retry-policy-ref")); + verify(mockElement, times(1)).getAttribute(eq("ssl-options-ref")); + verify(mockElement, times(1)).getAttribute(eq("contact-points")); + verify(mockElement, times(1)).getAttribute(eq("compression")); + verify(mockElement, times(1)).getAttribute(eq("jmx-reporting-enabled")); + verify(mockElement, times(1)).getAttribute(eq("metrics-enabled")); + verify(mockElement, times(1)).getAttribute(eq("password")); + verify(mockElement, times(1)).getAttribute(eq("port")); + verify(mockElement, times(1)).getAttribute(eq("ssl-enabled")); + verify(mockElement, times(1)).getAttribute(eq("username")); + } + + @Test + public void parseChildElementsWithLocalPoolingOptions() { + Element localPoolingOptionsElement = mock(Element.class); + + NodeList mockNodeList = mockNodeList(localPoolingOptionsElement); + + when(localPoolingOptionsElement.getLocalName()).thenReturn("local-pooling-options"); + when(localPoolingOptionsElement.getAttribute(eq("core-connections"))).thenReturn("50"); + when(localPoolingOptionsElement.getAttribute(eq("max-connections"))).thenReturn("200"); + when(localPoolingOptionsElement.getAttribute(eq("max-simultaneous-requests"))).thenReturn("50"); + when(localPoolingOptionsElement.getAttribute(eq("min-simultaneous-requests"))).thenReturn("5"); + when(mockElement.getChildNodes()).thenReturn(mockNodeList); + when(mockElement.getAttribute(eq("heartbeat-interval-seconds"))).thenReturn("15"); + when(mockElement.getAttribute(eq("idle-timeout-seconds"))).thenReturn("120"); + when(mockElement.getAttribute(eq("initialization-executor-ref"))).thenReturn("testExecutor"); + when(mockElement.getAttribute(eq("pool-timeout-milliseconds"))).thenReturn("60000"); + + BeanDefinitionBuilder builder = BeanDefinitionBuilder.genericBeanDefinition(); + + parser.parseChildElements(mockElement, mockParserContext(), builder); + + BeanDefinition beanDefinition = builder.getBeanDefinition(); + + BeanDefinition poolingOptionsBeanDefinition = getPropertyValue(beanDefinition, "poolingOptions"); + + assertThat(poolingOptionsBeanDefinition, is(notNullValue(BeanDefinition.class))); + assertThat(poolingOptionsBeanDefinition.getBeanClassName(), is(equalTo(PoolingOptionsFactoryBean.class.getName()))); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "heartbeatIntervalSeconds"), is(equalTo("15"))); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "idleTimeoutSeconds"), is(equalTo("120"))); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "initializationExecutor"), is(equalTo("testExecutor"))); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "poolTimeoutMilliseconds"), is(equalTo("60000"))); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "localCoreConnections"), is(equalTo("50"))); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "localMaxConnections"), is(equalTo("200"))); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "localMaxSimultaneousRequests"), is(equalTo("50"))); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "localMinSimultaneousRequests"), is(equalTo("5"))); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "remoteCoreConnections"), is(nullValue())); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "remoteMaxConnections"), is(nullValue())); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "remoteMaxSimultaneousRequests"), is(nullValue())); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "remoteMinSimultaneousRequests"), is(nullValue())); + + verify(mockElement, times(1)).getChildNodes(); + verify(mockElement, times(1)).getAttribute(eq("heartbeat-interval-seconds")); + verify(mockElement, times(1)).getAttribute(eq("idle-timeout-seconds")); + verify(mockElement, times(1)).getAttribute(eq("initialization-executor-ref")); + verify(mockElement, times(1)).getAttribute(eq("pool-timeout-milliseconds")); + verify(localPoolingOptionsElement, times(1)).getLocalName(); + verify(localPoolingOptionsElement, times(1)).getAttribute(eq("core-connections")); + verify(localPoolingOptionsElement, times(1)).getAttribute(eq("max-connections")); + verify(localPoolingOptionsElement, times(1)).getAttribute(eq("max-simultaneous-requests")); + verify(localPoolingOptionsElement, times(1)).getAttribute(eq("min-simultaneous-requests")); + } + + @Test + public void parseChildElementsWithRemotePoolingOptions() { + Element localPoolingOptionsElement = mock(Element.class); + + NodeList mockNodeList = mockNodeList(localPoolingOptionsElement); + + when(localPoolingOptionsElement.getLocalName()).thenReturn("remote-pooling-options"); + when(localPoolingOptionsElement.getAttribute(eq("core-connections"))).thenReturn("50"); + when(localPoolingOptionsElement.getAttribute(eq("max-connections"))).thenReturn("200"); + when(localPoolingOptionsElement.getAttribute(eq("max-simultaneous-requests"))).thenReturn("50"); + when(localPoolingOptionsElement.getAttribute(eq("min-simultaneous-requests"))).thenReturn("5"); + when(mockElement.getChildNodes()).thenReturn(mockNodeList); + when(mockElement.getAttribute(eq("heartbeat-interval-seconds"))).thenReturn("15"); + when(mockElement.getAttribute(eq("idle-timeout-seconds"))).thenReturn("120"); + when(mockElement.getAttribute(eq("initialization-executor-ref"))).thenReturn("testExecutor"); + when(mockElement.getAttribute(eq("pool-timeout-milliseconds"))).thenReturn("60000"); + + BeanDefinitionBuilder builder = BeanDefinitionBuilder.genericBeanDefinition(); + + parser.parseChildElements(mockElement, mockParserContext(), builder); + + BeanDefinition beanDefinition = builder.getBeanDefinition(); + + BeanDefinition poolingOptionsBeanDefinition = getPropertyValue(beanDefinition, "poolingOptions"); + + assertThat(poolingOptionsBeanDefinition, is(notNullValue(BeanDefinition.class))); + assertThat(poolingOptionsBeanDefinition.getBeanClassName(), is(equalTo(PoolingOptionsFactoryBean.class.getName()))); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "heartbeatIntervalSeconds"), is(equalTo("15"))); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "idleTimeoutSeconds"), is(equalTo("120"))); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "initializationExecutor"), is(equalTo("testExecutor"))); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "poolTimeoutMilliseconds"), is(equalTo("60000"))); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "localCoreConnections"), is(nullValue())); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "localMaxConnections"), is(nullValue())); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "localMaxSimultaneousRequests"), is(nullValue())); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "localMinSimultaneousRequests"), is(nullValue())); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "remoteCoreConnections"), is(equalTo("50"))); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "remoteMaxConnections"), is(equalTo("200"))); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "remoteMaxSimultaneousRequests"), is(equalTo( + "50"))); + assertThat(getPropertyValueAsString(poolingOptionsBeanDefinition, "remoteMinSimultaneousRequests"), is(equalTo( + "5"))); + + verify(mockElement, times(1)).getChildNodes(); + verify(mockElement, times(1)).getAttribute(eq("heartbeat-interval-seconds")); + verify(mockElement, times(1)).getAttribute(eq("idle-timeout-seconds")); + verify(mockElement, times(1)).getAttribute(eq("initialization-executor-ref")); + verify(mockElement, times(1)).getAttribute(eq("pool-timeout-milliseconds")); + verify(localPoolingOptionsElement, times(1)).getLocalName(); + verify(localPoolingOptionsElement, times(1)).getAttribute(eq("core-connections")); + verify(localPoolingOptionsElement, times(1)).getAttribute(eq("max-connections")); + verify(localPoolingOptionsElement, times(1)).getAttribute(eq("max-simultaneous-requests")); + verify(localPoolingOptionsElement, times(1)).getAttribute(eq("min-simultaneous-requests")); + } + + @Test + public void parseLocalPoolingOptionsProperlyConfiguresBeanDefinition() { + when(mockElement.getAttribute(eq("core-connections"))).thenReturn("50"); + when(mockElement.getAttribute(eq("max-connections"))).thenReturn("200"); + when(mockElement.getAttribute(eq("max-simultaneous-requests"))).thenReturn("50"); + when(mockElement.getAttribute(eq("min-simultaneous-requests"))).thenReturn("5"); + + BeanDefinitionBuilder builder = BeanDefinitionBuilder.genericBeanDefinition(); + + parser.parseLocalPoolingOptions(mockElement, builder); + + BeanDefinition beanDefinition = builder.getBeanDefinition(); + + assertThat(getPropertyValueAsString(beanDefinition, "heartbeatIntervalSeconds"), is(nullValue())); + assertThat(getPropertyValueAsString(beanDefinition, "idleTimeoutSeconds"), is(nullValue())); + assertThat(getPropertyValueAsString(beanDefinition, "initializationExecutor"), is(nullValue())); + assertThat(getPropertyValueAsString(beanDefinition, "poolTimeoutMilliseconds"), is(nullValue())); + assertThat(getPropertyValueAsString(beanDefinition, "localCoreConnections"), is(equalTo("50"))); + assertThat(getPropertyValueAsString(beanDefinition, "localMaxConnections"), is(equalTo("200"))); + assertThat(getPropertyValueAsString(beanDefinition, "localMaxSimultaneousRequests"), is(equalTo("50"))); + assertThat(getPropertyValueAsString(beanDefinition, "localMinSimultaneousRequests"), is(equalTo("5"))); + assertThat(getPropertyValueAsString(beanDefinition, "remoteCoreConnections"), is(nullValue())); + assertThat(getPropertyValueAsString(beanDefinition, "remoteMaxConnections"), is(nullValue())); + assertThat(getPropertyValueAsString(beanDefinition, "remoteMaxSimultaneousRequests"), is(nullValue())); + assertThat(getPropertyValueAsString(beanDefinition, "remoteMinSimultaneousRequests"), is(nullValue())); + + verify(mockElement, never()).getAttribute(eq("heartbeat-interval-seconds")); + verify(mockElement, never()).getAttribute(eq("idle-timeout-seconds")); + verify(mockElement, never()).getAttribute(eq("initialization-executor-ref")); + verify(mockElement, never()).getAttribute(eq("pool-timeout-milliseconds")); + verify(mockElement, times(1)).getAttribute(eq("core-connections")); + verify(mockElement, times(1)).getAttribute(eq("max-connections")); + verify(mockElement, times(1)).getAttribute(eq("max-simultaneous-requests")); + verify(mockElement, times(1)).getAttribute(eq("min-simultaneous-requests")); + } + + @Test + public void parseRemotePoolingOptionsProperlyConfiguresBeanDefinition() { + when(mockElement.getAttribute(eq("core-connections"))).thenReturn("50"); + when(mockElement.getAttribute(eq("max-connections"))).thenReturn("200"); + when(mockElement.getAttribute(eq("max-simultaneous-requests"))).thenReturn("50"); + when(mockElement.getAttribute(eq("min-simultaneous-requests"))).thenReturn("5"); + + BeanDefinitionBuilder builder = BeanDefinitionBuilder.genericBeanDefinition(); + + parser.parseRemotePoolingOptions(mockElement, builder); + + BeanDefinition beanDefinition = builder.getBeanDefinition(); + + assertThat(getPropertyValueAsString(beanDefinition, "heartbeatIntervalSeconds"), is(nullValue())); + assertThat(getPropertyValueAsString(beanDefinition, "idleTimeoutSeconds"), is(nullValue())); + assertThat(getPropertyValueAsString(beanDefinition, "initializationExecutor"), is(nullValue())); + assertThat(getPropertyValueAsString(beanDefinition, "poolTimeoutMilliseconds"), is(nullValue())); + assertThat(getPropertyValueAsString(beanDefinition, "localCoreConnections"), is(nullValue())); + assertThat(getPropertyValueAsString(beanDefinition, "localMaxConnections"), is(nullValue())); + assertThat(getPropertyValueAsString(beanDefinition, "localMaxSimultaneousRequests"), is(nullValue())); + assertThat(getPropertyValueAsString(beanDefinition, "localMinSimultaneousRequests"), is(nullValue())); + assertThat(getPropertyValueAsString(beanDefinition, "remoteCoreConnections"), is(equalTo("50"))); + assertThat(getPropertyValueAsString(beanDefinition, "remoteMaxConnections"), is(equalTo("200"))); + assertThat(getPropertyValueAsString(beanDefinition, "remoteMaxSimultaneousRequests"), is(equalTo("50"))); + assertThat(getPropertyValueAsString(beanDefinition, "remoteMinSimultaneousRequests"), is(equalTo("5"))); + + verify(mockElement, never()).getAttribute(eq("heartbeat-interval-seconds")); + verify(mockElement, never()).getAttribute(eq("idle-timeout-seconds")); + verify(mockElement, never()).getAttribute(eq("initialization-executor-ref")); + verify(mockElement, never()).getAttribute(eq("pool-timeout-milliseconds")); + verify(mockElement, times(1)).getAttribute(eq("core-connections")); + verify(mockElement, times(1)).getAttribute(eq("max-connections")); + verify(mockElement, times(1)).getAttribute(eq("max-simultaneous-requests")); + verify(mockElement, times(1)).getAttribute(eq("min-simultaneous-requests")); + } + + @Test + public void parseScript() { + when(mockElement.getTextContent()).thenReturn("CREATE TABLE schema.table;"); + assertThat(parser.parseScript(mockElement), is(equalTo("CREATE TABLE schema.table;"))); + verify(mockElement, times(1)).getTextContent(); + } + + @Test + public void newSocketOptionsBeanDefinitionIsProperlyInitialized() { + when(mockElement.getAttribute(eq("connect-timeout-millis"))).thenReturn("15000"); + when(mockElement.getAttribute(eq("keep-alive"))).thenReturn("true"); + when(mockElement.getAttribute(eq("read-timeout-millis"))).thenReturn("20000"); + when(mockElement.getAttribute(eq("receive-buffer-size"))).thenReturn("32768"); + when(mockElement.getAttribute(eq("reuse-address"))).thenReturn("true"); + when(mockElement.getAttribute(eq("send-buffer-size"))).thenReturn("16384"); + when(mockElement.getAttribute(eq("so-linger"))).thenReturn("false"); + when(mockElement.getAttribute(eq("tcp-no-delay"))).thenReturn("true"); + + BeanDefinition beanDefinition = parser.newSocketOptionsBeanDefinition(mockElement, mockParserContext()); + + assertThat(beanDefinition, is(notNullValue(BeanDefinition.class))); + assertThat(beanDefinition.getBeanClassName(), is(equalTo(SocketOptionsFactoryBean.class.getName()))); + assertThat((Element) beanDefinition.getSource(), is(equalTo(mockElement))); + assertThat(getPropertyValueAsString(beanDefinition, "connectTimeoutMillis"), is(equalTo("15000"))); + assertThat(getPropertyValueAsString(beanDefinition, "keepAlive"), is(equalTo("true"))); + assertThat(getPropertyValueAsString(beanDefinition, "readTimeoutMillis"), is(equalTo("20000"))); + assertThat(getPropertyValueAsString(beanDefinition, "receiveBufferSize"), is(equalTo("32768"))); + assertThat(getPropertyValueAsString(beanDefinition, "reuseAddress"), is(equalTo("true"))); + assertThat(getPropertyValueAsString(beanDefinition, "sendBufferSize"), is(equalTo("16384"))); + assertThat(getPropertyValueAsString(beanDefinition, "soLinger"), is(equalTo("false"))); + assertThat(getPropertyValueAsString(beanDefinition, "tcpNoDelay"), is(equalTo("true"))); + + verify(mockElement, times(1)).getAttribute(eq("connect-timeout-millis")); + verify(mockElement, times(1)).getAttribute(eq("keep-alive")); + verify(mockElement, times(1)).getAttribute(eq("read-timeout-millis")); + verify(mockElement, times(1)).getAttribute(eq("receive-buffer-size")); + verify(mockElement, times(1)).getAttribute(eq("reuse-address")); + verify(mockElement, times(1)).getAttribute(eq("send-buffer-size")); + verify(mockElement, times(1)).getAttribute(eq("so-linger")); + verify(mockElement, times(1)).getAttribute(eq("tcp-no-delay")); + } +} diff --git a/spring-cql/src/test/java/org/springframework/cassandra/config/xml/ParsingUtilsUnitTests.java b/spring-cql/src/test/java/org/springframework/cassandra/config/xml/ParsingUtilsUnitTests.java new file mode 100644 index 000000000..e9a9dc02a --- /dev/null +++ b/spring-cql/src/test/java/org/springframework/cassandra/config/xml/ParsingUtilsUnitTests.java @@ -0,0 +1,152 @@ +/* + * Copyright 2013-2016 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.config.xml; + +import static org.hamcrest.Matchers.*; +import static org.junit.Assert.*; + +import org.junit.Rule; +import org.junit.Test; +import org.junit.rules.ExpectedException; +import org.springframework.beans.PropertyValue; +import org.springframework.beans.factory.config.BeanDefinition; +import org.springframework.beans.factory.config.RuntimeBeanReference; +import org.springframework.beans.factory.support.BeanDefinitionBuilder; + +/** + * Test suite of unit tests testing the contract and functionality of the {@link ParsingUtils} class. + * + * @author John Blum + * @see org.springframework.cassandra.config.xml.ParsingUtils + * @since 1.5.0 + */ +// TODO: add more tests! +public class ParsingUtilsUnitTests { + + @Rule + public ExpectedException exception = ExpectedException.none(); + + @SuppressWarnings("unchecked") + protected T getPropertyValue(BeanDefinition beanDefinition, String propertyName) { + PropertyValue propertyValue = beanDefinition.getPropertyValues().getPropertyValue(propertyName); + + return (T) (propertyValue != null ? propertyValue.getValue() : null); + } + + @Test + public void addOptionalReferencePropertyUsesDefault() { + BeanDefinitionBuilder builder = ParsingUtils.addProperty(BeanDefinitionBuilder.genericBeanDefinition(), + "referenceProperty", null, "defaultBeanReference", false, true); + + RuntimeBeanReference propertyValue = getPropertyValue(builder.getBeanDefinition(), "referenceProperty"); + + assertThat(propertyValue, is(notNullValue(RuntimeBeanReference.class))); + assertThat(propertyValue.getBeanName(), is(equalTo("defaultBeanReference"))); + } + + @Test + public void addOptionalReferencePropertyWithNoValueDoesReturnsWithoutAdding() { + BeanDefinitionBuilder builder = ParsingUtils.addProperty(BeanDefinitionBuilder.genericBeanDefinition(), + "referenceProperty", null, null, false, false); + + BeanDefinition beanDefinition = builder.getRawBeanDefinition(); + + assertThat(beanDefinition.getPropertyValues().contains("referenceProperty"), is(false)); + assertThat(beanDefinition.getPropertyValues().isEmpty(), is(true)); + } + + @Test + public void addOptionalValuePropertyUsesDefault() { + BeanDefinitionBuilder builder = ParsingUtils.addProperty(BeanDefinitionBuilder.genericBeanDefinition(), + "valueProperty", null, "defaultValue", false, false); + + String propertyValue = getPropertyValue(builder.getBeanDefinition(), "valueProperty"); + + assertThat(propertyValue, is(equalTo("defaultValue"))); + } + + @Test + public void addOptionalValuePropertyWithNoValueDoesReturnsWithoutAdding() { + BeanDefinitionBuilder builder = ParsingUtils.addProperty(BeanDefinitionBuilder.genericBeanDefinition(), + "valueProperty", null, null, false, false); + + BeanDefinition beanDefinition = builder.getRawBeanDefinition(); + + assertThat(beanDefinition.getPropertyValues().contains("valueProperty"), is(false)); + assertThat(beanDefinition.getPropertyValues().isEmpty(), is(true)); + } + + @Test + public void addRequiredReferencePropertyIsSuccessful() { + BeanDefinitionBuilder builder = ParsingUtils.addProperty(BeanDefinitionBuilder.genericBeanDefinition(), + "referenceProperty", "reference", null, true, true); + + RuntimeBeanReference propertyValue = getPropertyValue(builder.getBeanDefinition(), "referenceProperty"); + + assertThat(propertyValue, is(notNullValue(RuntimeBeanReference.class))); + assertThat(propertyValue.getBeanName(), is(equalTo("reference"))); + } + + @Test + public void addRequiredReferencePropertyWithNoReferenceFails() { + exception.expect(IllegalArgumentException.class); + exception.expectCause(is(nullValue(Throwable.class))); + exception.expectMessage("value required for property reference [referenceProperty] on class [null]"); + + ParsingUtils.addProperty(BeanDefinitionBuilder.genericBeanDefinition(), "referenceProperty", null, + "defaultReference", true, true); + } + + @Test + public void addRequiredValuePropertyIsSuccessful() { + BeanDefinitionBuilder builder = ParsingUtils.addProperty(BeanDefinitionBuilder.genericBeanDefinition(), + "valueProperty", "value", null, true, false); + + String propertyValue = getPropertyValue(builder.getBeanDefinition(), "valueProperty"); + + assertThat(propertyValue, is(equalTo("value"))); + } + + @Test + public void addRequiredValuePropertyWithNoValueFails() { + exception.expect(IllegalArgumentException.class); + exception.expectCause(is(nullValue(Throwable.class))); + exception.expectMessage("value required for property [valueProperty] on class [null]"); + + ParsingUtils.addProperty(BeanDefinitionBuilder.genericBeanDefinition(), "valueProperty", null, + "defaultValue", true, false); + } + + @Test + public void addPropertyThrowsIllegalArgumentExceptionForNullBuilder() { + exception.expect(IllegalArgumentException.class); + exception.expectCause(is(nullValue(Throwable.class))); + exception.expectMessage("BeanDefinitionBuilder must not be null"); + + ParsingUtils.addProperty(null, "propertyName", "value", "defaultValue", false, false); + } + + @Test + public void addPropertyThrowsIllegalArgumentExceptionForNullPropertyName() { + exception.expect(IllegalArgumentException.class); + exception.expectCause(is(nullValue(Throwable.class))); + exception.expectMessage("Property name must not be null"); + + ParsingUtils.addProperty(BeanDefinitionBuilder.genericBeanDefinition(), null, "value", "defaultValue", + false, true); + } +} diff --git a/spring-cql/src/test/java/org/springframework/cassandra/test/integration/config/xml/PropertyPlaceholderNamespaceCreatingXmlConfigIntegrationTests.java b/spring-cql/src/test/java/org/springframework/cassandra/test/integration/config/xml/PropertyPlaceholderNamespaceCreatingXmlConfigIntegrationTests.java index c13e273f0..49a9ba4eb 100644 --- a/spring-cql/src/test/java/org/springframework/cassandra/test/integration/config/xml/PropertyPlaceholderNamespaceCreatingXmlConfigIntegrationTests.java +++ b/spring-cql/src/test/java/org/springframework/cassandra/test/integration/config/xml/PropertyPlaceholderNamespaceCreatingXmlConfigIntegrationTests.java @@ -15,36 +15,45 @@ */ package org.springframework.cassandra.test.integration.config.xml; +import static org.hamcrest.Matchers.*; import static org.junit.Assert.*; -import static org.springframework.cassandra.core.keyspace.DropKeyspaceSpecification.*; -import javax.inject.Inject; - -import org.junit.After; import org.junit.Test; import org.junit.runner.RunWith; +import org.springframework.beans.factory.annotation.Autowired; import org.springframework.cassandra.core.CqlOperations; import org.springframework.cassandra.test.integration.AbstractEmbeddedCassandraIntegrationTest; import org.springframework.cassandra.test.integration.config.IntegrationTestUtils; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; +import com.datastax.driver.core.Cluster; +import com.datastax.driver.core.HostDistance; +import com.datastax.driver.core.PoolingOptions; import com.datastax.driver.core.Session; +import com.datastax.driver.core.SocketOptions; /** * @author Mark Paluch + * @author John Blum */ @RunWith(SpringJUnit4ClassRunner.class) @ContextConfiguration +@SuppressWarnings("unused") public class PropertyPlaceholderNamespaceCreatingXmlConfigIntegrationTests extends AbstractEmbeddedCassandraIntegrationTest { - @Inject private Session session; + @Autowired + private Cluster cassandraCluster; - @Inject private CqlOperations ops; + @Autowired + private CqlOperations ops; + + @Autowired + private Session session; @Test - public void test() { + public void keyspaceExists() { IntegrationTestUtils.assertSession(session); IntegrationTestUtils.assertKeyspaceExists("ppncxct", session); @@ -52,9 +61,38 @@ public class PropertyPlaceholderNamespaceCreatingXmlConfigIntegrationTests assertNotNull(ops); } - @After - public void tearDown() throws Exception { - dropKeyspace("ppncxct"); - dropKeyspace("foo123"); + @Test + public void localAndRemotePoolingOptionsWereConfiguredProperly() { + + PoolingOptions poolingOptions = cassandraCluster.getConfiguration().getPoolingOptions(); + + assertThat(poolingOptions, is(notNullValue(PoolingOptions.class))); + assertThat(poolingOptions.getHeartbeatIntervalSeconds(), is(equalTo(60))); + assertThat(poolingOptions.getIdleTimeoutSeconds(), is(equalTo(180))); + assertThat(poolingOptions.getPoolTimeoutMillis(), is(equalTo(30000))); + assertThat(poolingOptions.getCoreConnectionsPerHost(HostDistance.LOCAL), is(equalTo(4))); + assertThat(poolingOptions.getMaxConnectionsPerHost(HostDistance.LOCAL), is(equalTo(8))); + assertThat(poolingOptions.getMaxRequestsPerConnection(HostDistance.LOCAL), is(equalTo(20))); + assertThat(poolingOptions.getNewConnectionThreshold(HostDistance.LOCAL), is(equalTo(10))); + assertThat(poolingOptions.getCoreConnectionsPerHost(HostDistance.REMOTE), is(equalTo(2))); + assertThat(poolingOptions.getMaxConnectionsPerHost(HostDistance.REMOTE), is(equalTo(4))); + assertThat(poolingOptions.getMaxRequestsPerConnection(HostDistance.REMOTE), is(equalTo(10))); + assertThat(poolingOptions.getNewConnectionThreshold(HostDistance.REMOTE), is(equalTo(5))); + } + + @Test + public void socketOptionsWereConfiguredProperly() { + + SocketOptions socketOptions = cassandraCluster.getConfiguration().getSocketOptions(); + + assertThat(socketOptions, is(notNullValue(SocketOptions.class))); + assertThat(socketOptions.getConnectTimeoutMillis(), is(equalTo(15000))); + assertThat(socketOptions.getKeepAlive(), is(true)); + assertThat(socketOptions.getReadTimeoutMillis(), is(equalTo(60000))); + assertThat(socketOptions.getReceiveBufferSize(), is(equalTo(1024))); + assertThat(socketOptions.getReuseAddress(), is(true)); + assertThat(socketOptions.getSendBufferSize(), is(equalTo(2048))); + assertThat(socketOptions.getSoLinger(), is(equalTo(5))); + assertThat(socketOptions.getTcpNoDelay(), is(false)); } } diff --git a/spring-cql/src/test/java/org/springframework/cassandra/test/integration/config/xml/XmlConfigIntegrationTests.java b/spring-cql/src/test/java/org/springframework/cassandra/test/integration/config/xml/XmlConfigIntegrationTests.java index 3fec09d12..e156f4ee6 100644 --- a/spring-cql/src/test/java/org/springframework/cassandra/test/integration/config/xml/XmlConfigIntegrationTests.java +++ b/spring-cql/src/test/java/org/springframework/cassandra/test/integration/config/xml/XmlConfigIntegrationTests.java @@ -15,6 +15,11 @@ */ package org.springframework.cassandra.test.integration.config.xml; +import static org.hamcrest.Matchers.*; +import static org.junit.Assert.*; + +import java.util.concurrent.Executor; + import org.junit.After; import org.junit.Before; import org.junit.Rule; @@ -25,36 +30,91 @@ import org.springframework.cassandra.test.integration.config.IntegrationTestUtil import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.support.ClassPathXmlApplicationContext; +import com.datastax.driver.core.Cluster; +import com.datastax.driver.core.HostDistance; +import com.datastax.driver.core.PoolingOptions; import com.datastax.driver.core.Session; +import com.datastax.driver.core.SocketOptions; /** + * Test XML namespace configuration using the spring-cql-1.0.xsd. + * * @author Matthews T. Adams * @author Oliver Gierke * @author Mark Paluch + * @author John Blum */ +@SuppressWarnings("unused") public class XmlConfigIntegrationTests extends AbstractEmbeddedCassandraIntegrationTest { public static final String KEYSPACE = "xmlconfigtest"; - @Rule public KeyspaceRule keyspaceRule = new KeyspaceRule(cassandraEnvironment, KEYSPACE); + @Rule + public KeyspaceRule keyspaceRule = new KeyspaceRule(cassandraEnvironment, KEYSPACE); + + private ConfigurableApplicationContext applicationContext; + + private Cluster cluster; + + private Executor executor; private Session session; - private ConfigurableApplicationContext context; @Before public void setUp() { + this.applicationContext = new ClassPathXmlApplicationContext( + "XmlConfigIntegrationTests-context.xml", getClass()); - this.context = new ClassPathXmlApplicationContext("XmlConfigIntegrationTests-context.xml", getClass()); - this.session = context.getBean(Session.class); + this.cluster = applicationContext.getBean(Cluster.class); + this.executor = applicationContext.getBean(Executor.class); + this.session = applicationContext.getBean(Session.class); } @After public void tearDown() { - context.close(); + if (this.applicationContext != null) { + this.applicationContext.close(); + } } @Test - public void test() { + public void keyspaceExists() { IntegrationTestUtils.assertKeyspaceExists(KEYSPACE, session); } + + @Test + public void localAndRemotePoolingOptionsWereConfiguredProperly() { + + PoolingOptions poolingOptions = cluster.getConfiguration().getPoolingOptions(); + + assertThat(poolingOptions, is(notNullValue(PoolingOptions.class))); + assertThat(poolingOptions.getHeartbeatIntervalSeconds(), is(equalTo(60))); + assertThat(poolingOptions.getIdleTimeoutSeconds(), is(equalTo(300))); + assertThat(poolingOptions.getInitializationExecutor(), is(equalTo(executor))); + assertThat(poolingOptions.getPoolTimeoutMillis(), is(equalTo(15000))); + assertThat(poolingOptions.getCoreConnectionsPerHost(HostDistance.LOCAL), is(equalTo(2))); + assertThat(poolingOptions.getMaxConnectionsPerHost(HostDistance.LOCAL), is(equalTo(8))); + assertThat(poolingOptions.getMaxRequestsPerConnection(HostDistance.LOCAL), is(equalTo(100))); + assertThat(poolingOptions.getNewConnectionThreshold(HostDistance.LOCAL), is(equalTo(25))); + assertThat(poolingOptions.getCoreConnectionsPerHost(HostDistance.REMOTE), is(equalTo(1))); + assertThat(poolingOptions.getMaxConnectionsPerHost(HostDistance.REMOTE), is(equalTo(2))); + assertThat(poolingOptions.getMaxRequestsPerConnection(HostDistance.REMOTE), is(equalTo(100))); + assertThat(poolingOptions.getNewConnectionThreshold(HostDistance.REMOTE), is(equalTo(25))); + } + + @Test + public void socketOptionsWereConfiguredProperly() { + + SocketOptions socketOptions = cluster.getConfiguration().getSocketOptions(); + + assertThat(socketOptions, is(notNullValue(SocketOptions.class))); + assertThat(socketOptions.getConnectTimeoutMillis(), is(equalTo(5000))); + assertThat(socketOptions.getKeepAlive(), is(true)); + assertThat(socketOptions.getReadTimeoutMillis(), is(equalTo(60000))); + assertThat(socketOptions.getReceiveBufferSize(), is(equalTo(65536))); + assertThat(socketOptions.getReuseAddress(), is(true)); + assertThat(socketOptions.getSendBufferSize(), is(equalTo(65536))); + assertThat(socketOptions.getSoLinger(), is(equalTo(60))); + assertThat(socketOptions.getTcpNoDelay(), is(true)); + } } diff --git a/spring-cql/src/test/resources/org/springframework/cassandra/test/integration/config/xml/PropertyPlaceholderNamespaceCreatingXmlConfigIntegrationTests-context.xml b/spring-cql/src/test/resources/org/springframework/cassandra/test/integration/config/xml/PropertyPlaceholderNamespaceCreatingXmlConfigIntegrationTests-context.xml index f15ff72f3..7da0dbdc0 100644 --- a/spring-cql/src/test/resources/org/springframework/cassandra/test/integration/config/xml/PropertyPlaceholderNamespaceCreatingXmlConfigIntegrationTests-context.xml +++ b/spring-cql/src/test/resources/org/springframework/cassandra/test/integration/config/xml/PropertyPlaceholderNamespaceCreatingXmlConfigIntegrationTests-context.xml @@ -1,13 +1,16 @@ + xmlns:cassandra="http://www.springframework.org/schema/cql" + xmlns:context="http://www.springframework.org/schema/context" + xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" + xsi:schemaLocation=" + http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd + http://www.springframework.org/schema/context http://www.springframework.org/schema/context/spring-context.xsd + http://www.springframework.org/schema/cql http://www.springframework.org/schema/cql/spring-cql.xsd +"> + location="classpath:/org/springframework/cassandra/test/integration/config/xml/ppncxct.properties"/> @@ -15,66 +18,72 @@ + value="${cluster.reconnection.delayMillis}"/> + class="com.datastax.driver.core.policies.DowngradingConsistencyRetryPolicy"/> + class="org.springframework.cassandra.test.integration.config.xml.TestHostStateListener"/> + class="org.springframework.cassandra.test.integration.config.xml.TestLatencyTracker"/> - + + core-connections="${local.core.connections}" + max-connections="${local.max.connections}" + max-simultaneous-requests="${local.max.requests}" + min-simultaneous-requests="${local.min.requests}"/> + core-connections="${remote.core.connections}" + max-connections="${remote.max.connections}" + max-simultaneous-requests="${remote.max.requests}" + min-simultaneous-requests="${remote.min.requests}"/> + connect-timeout-millis="${socket.connectTimeoutMillis}" + keep-alive="${socket.keepAlive}" + read-timeout-millis="${socket.readTimeoutMillis}" + receive-buffer-size="${socket.receiveBufferSize}" + reuse-address="${socket.reuseAddress}" + send-buffer-size="${socket.sendBufferSize}" + so-linger="${socket.soLinger}" + tcp-no-delay="${socket.tcpNoDelay}"/> - + - - + + - + - + diff --git a/spring-cql/src/test/resources/org/springframework/cassandra/test/integration/config/xml/XmlConfigIntegrationTests-context.xml b/spring-cql/src/test/resources/org/springframework/cassandra/test/integration/config/xml/XmlConfigIntegrationTests-context.xml index 58b3b398a..35c3fc561 100644 --- a/spring-cql/src/test/resources/org/springframework/cassandra/test/integration/config/xml/XmlConfigIntegrationTests-context.xml +++ b/spring-cql/src/test/resources/org/springframework/cassandra/test/integration/config/xml/XmlConfigIntegrationTests-context.xml @@ -1,29 +1,46 @@ + xmlns:cassandra="http://www.springframework.org/schema/cql" + xmlns:context="http://www.springframework.org/schema/context" + xmlns:task="http://www.springframework.org/schema/task" + xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" + xsi:schemaLocation=" + http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd + http://www.springframework.org/schema/context http://www.springframework.org/schema/context/spring-context.xsd + http://www.springframework.org/schema/cql http://www.springframework.org/schema/cql/spring-cql.xsd + http://www.springframework.org/schema/task http://www.springframework.org/schema/task/spring-task.xsd +"> - + - - + + + + + core-connections="1" + max-connections="2" + max-simultaneous-requests="100" + min-simultaneous-requests="25"/> + connect-timeout-millis="5000" + keep-alive="true" + read-timeout-millis="60000" + receive-buffer-size="65536" + reuse-address="true" + send-buffer-size="65536" + so-linger="60" + tcp-no-delay="true"/> - + diff --git a/spring-cql/src/test/resources/org/springframework/cassandra/test/integration/config/xml/ppncxct.properties b/spring-cql/src/test/resources/org/springframework/cassandra/test/integration/config/xml/ppncxct.properties index ee21af18d..a9bdaa517 100644 --- a/spring-cql/src/test/resources/org/springframework/cassandra/test/integration/config/xml/ppncxct.properties +++ b/spring-cql/src/test/resources/org/springframework/cassandra/test/integration/config/xml/ppncxct.properties @@ -5,6 +5,9 @@ cluster.metricsEnabled=false cluster.jmxReportingEnabled=false cluster.reconnection.delayMillis=5000 cluster.sslEnabled= true +cluster.poolingoptions.heartbeatintervalseconds=60 +cluster.poolingoptions.idletimeoutseconds=180 +cluster.poolingoptions.pooltimeoutmilliseconds=30000 keyspace.name=ppncxct keyspace.action=CREATE_DROP dc1.name=DCJAX @@ -15,19 +18,19 @@ lb.policy.dcAware.remoteHosts=0 lb.policy.dcAware.localDc=DCJAX auth.username=test auth.password=pass -socket.connectTimeoutMillis=5000 +socket.connectTimeoutMillis=15000 socket.keepAlive=true socket.readTimeoutMillis=60000 socket.receiveBufferSize=1024 -socket.sendBufferSize=2048 socket.reuseAddress=true +socket.sendBufferSize=2048 socket.soLinger=5 socket.tcpNoDelay=false -local.min.requests=10 -local.max.requests=20 local.core.connections=4 local.max.connections=8 -remote.min.requests=5 -remote.max.requests=10 +local.max.requests=20 +local.min.requests=10 remote.core.connections=2 remote.max.connections=4 +remote.max.requests=10 +remote.min.requests=5 diff --git a/spring-data-cassandra/src/main/resources/org/springframework/data/cassandra/config/spring-cassandra-1.0.xsd b/spring-data-cassandra/src/main/resources/org/springframework/data/cassandra/config/spring-cassandra-1.0.xsd index 448702c70..7ff9ed16c 100644 --- a/spring-data-cassandra/src/main/resources/org/springframework/data/cassandra/config/spring-cassandra-1.0.xsd +++ b/spring-data-cassandra/src/main/resources/org/springframework/data/cassandra/config/spring-cassandra-1.0.xsd @@ -1,12 +1,13 @@ - + xmlns:beans="http://www.springframework.org/schema/beans" + xmlns:context="http://www.springframework.org/schema/context" + xmlns:repository="http://www.springframework.org/schema/data/repository" + xmlns:tool="http://www.springframework.org/schema/tool" + xmlns:xsd="http://www.w3.org/2001/XMLSchema" + xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" + targetNamespace="http://www.springframework.org/schema/data/cassandra" + elementFormDefault="qualified" attributeFormDefault="unqualified"> @@ -19,38 +20,20 @@ - + - - - + + - - - - - - - - - - - - + + @@ -66,42 +49,48 @@ Defines a Cassandra cluster. + + + + + + + + + + + - + - + - + - + - + - + - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + - + - - - - - - - - - - - - - - - - - - - - - - - - - - + + + - - + + - - + @@ -334,21 +356,162 @@ Custom Latency Tracker for the Cassandra Cluster. + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + - - - - - - - - - - + + + + + + + + + + @@ -387,6 +564,13 @@ Sets read timeout for client socket in milliseconds. ]]> + + + + + + + + + + - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - @@ -558,8 +608,7 @@ Provides the ability to configure the keyspace's replication settings. ]]> - + - + - + - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - @@ -713,6 +643,28 @@ The reference to a cassandraTemplate. Will default to 'cassandraTemplate'. + + + + + + + + + + + + + + + + + + + - + - + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + - + - + - + - + +