SGF-550 - Added GatewaySenders and GatewaySender annotation.

This commit is contained in:
kohlmu-pivotal
2019-03-26 09:38:45 -07:00
committed by John Blum
parent 7e99457c39
commit a6d86147c5
21 changed files with 2888 additions and 83 deletions

View File

@@ -225,7 +225,6 @@
<version>${spring-shell.version}</version>
<scope>test</scope>
</dependency>
</dependencies>
<build>

View File

@@ -650,6 +650,55 @@ To enable a Gateway Receiver the application class needs to be annotated with `@
}
}
----
NOTE: {data-store-name} is a server-side feature only and can only be configured on a CacheServer or PeerServer
NOTE: {data-store-name} GatewayReceiver is a server-side feature only and can only be configured on a CacheServer or PeerServer
See {sdg-javadoc}/org/springframework/data/gemfire/wan/annotation/EnableGatewayReceiver.html[`@EnableGatewayReceiver` Javadoc].
See {sdg-javadoc}/org/springframework/data/gemfire/wan/annotation/EnableGatewayReceiver.html[`@EnableGatewayReceiver` Javadoc].
[[bootstap-annotations-quickstart-gatewaysenders]]
To enable Gateway Senders in the application class needs to be annotated with `@EnableGatewaySenders` and `@EnableGatewaySender` as follows:
[source,java]
----
@Configuration
@CacheServerApplication
@EnableGatewaySenders(gatewaySenders = {
@EnableGatewaySender(name = "GatewaySender", manualStart = true,
remoteDistributedSystemId = 2, diskSynchronous = true, batchConflationEnabled = true,
parallel = true, persistent = false,diskStoreReference = "someDiskStore",
orderPolicy = OrderPolicyType.PARTITION, alertThreshold = 1234, batchSize = 100,
eventFilters = "SomeEventFilter", batchTimeInterval = 2000, dispatcherThreads = 22,
maximumQueueMemory = 400,socketBufferSize = 16384,
socketReadTimeout = 4000, regions = { "Region1"}),
@EnableGatewaySender(name = "GatewaySender2", manualStart = true,
remoteDistributedSystemId = 2, diskSynchronous = true, batchConflationEnabled = true,
parallel = true, persistent = false, diskStoreReference = "someDiskStore",
orderPolicy = OrderPolicyType.PARTITION, alertThreshold = 1234, batchSize = 100,
eventFilters = "SomeEventFilter", batchTimeInterval = 2000, dispatcherThreads = 22,
maximumQueueMemory = 400, socketBufferSize = 16384,socketReadTimeout = 4000,
regions = { "Region2" })
}){
....
}
----
NOTE: {data-store-name} GatewaySender is a server-side feature only and can only be configured on a CacheServer or PeerServer
In the above example, the application is configured with 2 Regions, `Region1` and `Region2`. In addition 2 GatewaySenders will be configured for service both Regions. `GatewaySender1` will be configured to replicate `Region1`'s data and `GatewaySender2` will be configured to replicate `Region2`'s data. As demonstrated each GatewaySender property can be configured on each `EnableGatewaySender` annotation.
It is also possible to have a more generic, "defaulted" properties approach, where all properties are configured on the `EnableGatewaySenders` annotation. This way a set of generic, defaulted values can be set on the parent annotation and then overridden on the child if required, as demonstrated below:
[source,java]
----
@EnableGatewaySenders(gatewaySenders = {
@EnableGatewaySender(name = "GatewaySender", transportFilters = "transportBean1", regions = "Region2"),
@EnableGatewaySender(name = "GatewaySender2")},
manualStart = true, remoteDistributedSystemId = 2,
diskSynchronous = false, batchConflationEnabled = true, parallel = true, persistent = true,
diskStoreReference = "someDiskStore", orderPolicy = OrderPolicyType.PARTITION, alertThreshold = 1234, batchSize = 1002,
eventFilters = "SomeEventFilter", batchTimeInterval = 2000, dispatcherThreads = 22, maximumQueueMemory = 400,
socketBufferSize = 16384, socketReadTimeout = 4000, regions = { "Region1", "Region2" },
transportFilters = { "transportBean2", "transportBean1" })
----
NOTE: When the `regions` property is left empty or not populated, the GatewaySender(s) will automatically attach itself to every
configured Region within the application
See {sdg-javadoc}/org/springframework/data/gemfire/wan/annotation/EnableGatewaySenders.html[`@EnableGatewaySenders` Javadoc] and {sdg-javadoc}/org/springframework/data/gemfire/wan/annotation/EnableGatewaySender.html[`@EnableGatewaySender` Javadoc].

View File

@@ -342,7 +342,8 @@ configuration metadata at runtime, before the Spring managed beans that the anno
* `PeerCacheConfigurer`
* `PoolConfigurer`
* `RegionConfigurer`
* `EnableGatewayReceiverConfigurer`
* `GatewayReceiverConfigurer`
* `GatewaySenderConfigurer`
For example, you can use the `CacheServerConfigurer` and `ClientCacheConfigurer` to customize the port numbers
used by your Spring Boot `CacheServer` and `ClientCache` applications, respectively.

View File

@@ -0,0 +1,205 @@
package org.springframework.data.gemfire.config.annotation;
import java.lang.annotation.Documented;
import java.lang.annotation.ElementType;
import java.lang.annotation.Inherited;
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
import java.lang.annotation.Target;
import org.springframework.context.annotation.Import;
import org.springframework.data.gemfire.config.support.GatewaySenderBeanFactoryPostProcessor;
import org.springframework.data.gemfire.wan.OrderPolicyType;
/**
* This annotation is responsible for the configuration of a single {@link org.apache.geode.cache.wan.GatewaySender}.
* All properties configured on this annotation will be be overrides from the defaults set on the {@link EnableGatewaySenders}.
*
* @author Udo Kohlmeyer
* @see EnableGatewaySenders
* @see org.apache.geode.cache.wan.GatewaySender
* @see org.apache.geode.cache.wan.GatewayReceiver
* @see org.apache.geode.cache.wan.GatewayEventFilter
* @see org.apache.geode.cache.wan.GatewayTransportFilter
* @see org.apache.geode.cache.wan.GatewaySender.OrderPolicy
* @see org.apache.geode.cache.wan.GatewayEventSubstitutionFilter
* @since 2.2.0
*/
@Target(ElementType.TYPE)
@Retention(RetentionPolicy.RUNTIME)
@Inherited
@Documented
@Import({ GatewaySenderBeanFactoryPostProcessor.class, GatewaySenderConfiguration.class })
@SuppressWarnings("unused")
public @interface EnableGatewaySender {
/**
* This property configures the time, in milliseconds, that an object can be in the queue to be replicated before the
* {@link org.apache.geode.cache.wan.GatewaySender} logs an alert.
* <p> This property can also be configured using the {@literal spring.data.gemfire.gateway.sender.<name>.alert-threshold}
* <p>Default value is {@link GatewaySenderConfiguration#DEFAULT_ALERT_THRESHOLD}
*/
int alertThreshold() default GatewaySenderConfiguration.DEFAULT_ALERT_THRESHOLD;
/**
* A boolean flag to indicate if the configured GatewaySender(s) should use conflate entries in each batch. This means,
* that a batch will never contain duplicate entries, as the batch will always only contain the latest value for a key.
* <p> This property can also be configured using the {@literal spring.data.gemfire.gateway.sender.<name>.batch-conflation-enabled}
* <p>Default value is {@value GatewaySenderConfiguration#DEFAULT_BATCH_CONFLATION_ENABLED}
*/
boolean batchConflationEnabled() default GatewaySenderConfiguration.DEFAULT_BATCH_CONFLATION_ENABLED;
/**
* This property configures the maximum batch size that the {@link org.apache.geode.cache.wan.GatewaySender} send to the
* remote site. This property works in conjunction with the {@link EnableGatewaySenders#batchTimeInterval()} setting.
* The {@link org.apache.geode.cache.wan.GatewaySender} will send when either the {@literal batch-size} or {@literal batch-time-interval}
* is met.
* <p> This property can also be configured using the {@literal spring.data.gemfire.gateway.sender.<name>.batch-size}
* <p>Default value is {@value GatewaySenderConfiguration#DEFAULT_BATCH_SIZE}
*/
int batchSize() default GatewaySenderConfiguration.DEFAULT_BATCH_SIZE;
/**
* This property configures the maximum batch time interval, in milliseconds, that the {@link org.apache.geode.cache.wan.GatewaySender} wait before
* attempting to send a batch of queued object to the remote {@link org.apache.geode.cache.wan.GatewayReceiver}.
* <p>This property works in conjunction with the {@link EnableGatewaySenders#batchSize()} setting.
* The {@link org.apache.geode.cache.wan.GatewaySender} will send when either the {@literal batch-size} or {@literal batch-time-interval}
* is met.
* <p> This property can also be configured using the {@literal spring.data.gemfire.gateway.sender.<name>.batch-time-interval}
* <p>Default value is {@value GatewaySenderConfiguration#DEFAULT_BATCH_TIME_INTERVAL}
*/
int batchTimeInterval() default GatewaySenderConfiguration.DEFAULT_BATCH_TIME_INTERVAL;
/**
* This property configures what {@link org.apache.geode.cache.DiskStore} the GatewaySender(s) are to use when persisting
* the GatewaySender(s) queues. This setting should be set when the {@literal persistent} property is set to {@literal true}.
* <p> This property can also be configured using the {@literal spring.data.gemfire.gateway.sender.<name>.diskstore-reference}
* <p>Default value is {@value GatewaySenderConfiguration#DEFAULT_DISK_STORE_REFERENCE}
*/
String diskStoreReference() default GatewaySenderConfiguration.DEFAULT_DISK_STORE_REFERENCE;
/**
* A boolean flag to indicate if the configured GatewaySender(s) should use synchronous diskstore writes.
* <p> This property can also be configured using the {@literal spring.data.gemfire.gateway.sender.<name>.disk-synchronous}
* <p>Default value is {@value GatewaySenderConfiguration#DEFAULT_DISK_SYNCHRONOUS}
*/
boolean diskSynchronous() default GatewaySenderConfiguration.DEFAULT_DISK_SYNCHRONOUS;
/**
* This property configures the number of dispatcher threads that the {@link org.apache.geode.cache.wan.GatewaySender}
* will try to use to dispatch the queuing object.
* <p>This property can also be configured using the {@literal spring.data.gemfire.gateway.sender.<name>.dispatcher-threads}
* <p>Default value is {@value GatewaySenderConfiguration#DEFAULT_DISPATCHER_THREADS}
*/
int dispatcherThreads() default GatewaySenderConfiguration.DEFAULT_DISPATCHER_THREADS;
/**
* A list of {@link org.apache.geode.cache.wan.GatewayEventFilter} to be applied to the {@link org.apache.geode.cache.wan.GatewaySender}.
* {@link org.apache.geode.cache.wan.GatewayEventFilter} are used to filter out objects from the sending queue before dispatching
* them to the remote {@link org.apache.geode.cache.wan.GatewayReceiver}.
* <p>This property can also be configured using the {@literal spring.data.gemfire.gateway.sender.<name>.event-filters} property.
* <p>Default value is and empty list
*/
String[] eventFilters() default {};
/**
* This property configures the {@link org.apache.geode.cache.wan.GatewayEventSubstitutionFilter} to be used by the GatewaySender(s).
* The {@link org.apache.geode.cache.wan.GatewayEventSubstitutionFilter} is used to replace values on objects before they
* are enqueue for remote replication.
* <p> This property can also be configured using the {@literal spring.data.gemfire.gateway.sender.<name>.event-substitution-filter}
* <p>Default value is {@link GatewaySenderConfiguration#DEFAULT_EVENT_SUBSTITUTION_FILTER}
*/
String eventSubstitutionFilter() default GatewaySenderConfiguration.DEFAULT_EVENT_SUBSTITUTION_FILTER;
/**
* A boolean to indicate if a configured {@link org.apache.geode.cache.wan.GatewaySender} should be started automatically
* <p>This property can also be configured using the {@literal spring.data.gemfire.gateway.sender.<name>.manual-start}
* <p>Default is {@value @EnableGatewaySenderConfiguration.DEFAULT_MANUAL_START}
*/
boolean manualStart() default GatewaySenderConfiguration.DEFAULT_MANUAL_START;
/**
* This property configures the maximum size, in MB, that the {@link org.apache.geode.cache.wan.GatewaySender} that the queue
* may take on heap memory, before overflowing to disk.
* <p>This property can also be configured using the {@literal spring.data.gemfire.gateway.sender.<name>.maximum-queue-memory}
* <p>Default value is {@value GatewaySenderConfiguration#DEFAULT_MAXIMUM_QUEUE_MEMORY}
*/
int maximumQueueMemory() default GatewaySenderConfiguration.DEFAULT_MAXIMUM_QUEUE_MEMORY;
/**
* Specifies the {@link String name} of the {@link org.apache.geode.cache.wan.GatewaySender}.
* This name is also used as the name of the bean registered in the Spring container as well as the name used in the resolution
* {@link org.apache.geode.cache.wan.GatewaySender} properties from {@literal application.properties}
* (e.g. {@literal spring.data.gemfire.gateway.sender.<name>.manual-start}), that are specific to this {@link org.apache.geode.cache.wan.GatewaySender}
*/
String name();
/**
* This property sets the ordering policy that the GatewaySender(s) will use when queueing entries to be replicated to
* a remote {@link org.apache.geode.cache.wan.GatewayReceiver}.
* <p>There are three different ordering policies:
* <ul>
* <li>{@link OrderPolicyType#KEY} - Order of events preserved on a per key basis</li>
* <li>{@link OrderPolicyType#THREAD} - Order of events preserved by the thread that added the event</li>
* <li>{@link OrderPolicyType#PARTITION} - Order of events is preserved in order that they arrived in partitioned Region</li>
* </ul>
* <p>This property can also be configured using the {@literal spring.data.gemfire.gateway.sender.<name>.order-policy}</p>
* <p>Default value is an empty list
*/
OrderPolicyType orderPolicy() default OrderPolicyType.KEY;
/**
* A boolean flag to indicate if the configured GatewaySender(s) should use parallel GatewaySender replication.
* Parallel replication means that each {@link org.apache.geode.cache.server.CacheServer} that defines a {@link org.apache.geode.cache.wan.GatewaySender}
* will send data to a remote {@link org.apache.geode.cache.wan.GatewayReceiver}.
* <p> This property can also be configured using the {@literal spring.data.gemfire.gateway.sender.<name>.parallel}
* <p>Default value is {@value GatewaySenderConfiguration#DEFAULT_PARALLEL}
*/
boolean parallel() default GatewaySenderConfiguration.DEFAULT_PARALLEL;
/**
* A boolean flag to indicate if the configured GatewaySender(s) should use persistence. This setting should be used
* in conjunction with the {@literal disk-store-reference} property.
* <p> This property can also be configured using the {@literal spring.data.gemfire.gateway.sender.<name>.persistent}
* <p>Default value is {@value GatewaySenderConfiguration#DEFAULT_PERSISTENT}
*/
boolean persistent() default GatewaySenderConfiguration.DEFAULT_PERSISTENT;
/**
* A list of {@link org.apache.geode.cache.Region} names that are to be configured with {@link org.apache.geode.cache.wan.GatewaySender} replication.
* An empty list will denote that ALL regions are to be replicated to the remote {@link org.apache.geode.cache.wan.GatewayReceiver}.
* <p>This property can also be configured using the {@literal spring.data.gemfire.gateway.sender.<name>.region-names} property.
* <p>Default value is an empty list
*/
String[] regions() default {};
/**
* The id of the distributed system (cluster) that the GatewaySender(s) would send their data to.
* <p>This property can also be configured using the {@literal spring.data.gemfire.gateway.sender.<name>.remote-distributed-system-id}
* <p>Default is {@value @EnableGatewaySenderConfiguration.DEFAULT_REMOTE_DISTRIBUTED_SYSTEM_ID}
*/
int remoteDistributedSystemId() default GatewaySenderConfiguration.DEFAULT_REMOTE_DISTRIBUTED_SYSTEM_ID;
/**
* The socket buffer size for the {@link org.apache.geode.cache.wan.GatewaySender}. This setting is in bytes.
* <p>This property can also be configured using the {@literal spring.data.gemfire.gateway.sender.<name>.socket-buffer-size} property.
* <p>Default value is {@value GatewaySenderConfiguration#DEFAULT_SOCKET_BUFFER_SIZE}
*/
int socketBufferSize() default GatewaySenderConfiguration.DEFAULT_SOCKET_BUFFER_SIZE;
/**
* Amount of time in milliseconds that the gateway sender will wait to receive an acknowledgment from a remote site.
* By default this is set to 0, which means there is no timeout. The minimum allowed timeout is 30000 (milliseconds).
* <p>This property can also be configured using the {@literal spring.data.gemfire.gateway.sender.<name>.socket-read-timeout} property.
* <p>Default value is {@value GatewaySenderConfiguration#DEFAULT_SOCKET_READ_TIMEOUT}
*/
int socketReadTimeout() default GatewaySenderConfiguration.DEFAULT_SOCKET_READ_TIMEOUT;
/**
* An in-order list of {@link org.apache.geode.cache.wan.GatewayTransportFilter} to be applied to
* {@link org.apache.geode.cache.wan.GatewaySender}
* <p>This property can also be configured using the {@literal spring.data.gemfire.gateway.sender.<name>.transport-filters}s property
* <p>Default value is an empty list
*/
String[] transportFilters() default {};
}

View File

@@ -0,0 +1,202 @@
package org.springframework.data.gemfire.config.annotation;
import java.lang.annotation.Documented;
import java.lang.annotation.ElementType;
import java.lang.annotation.Inherited;
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
import java.lang.annotation.Target;
import org.springframework.context.annotation.Import;
import org.springframework.data.gemfire.config.support.GatewaySenderBeanFactoryPostProcessor;
import org.springframework.data.gemfire.wan.OrderPolicyType;
/**
* This annotation is responsible for the configuration of {@link org.apache.geode.cache.wan.GatewaySender}.
* All properties configured on this annotation will be used as default value for all {@link org.apache.geode.cache.wan.GatewaySender}
* configured within the <b><i>gatewaySenders</i></b> property.
*
* @author Udo Kohlmeyer
* @see org.apache.geode.cache.wan.GatewaySender
* @see org.apache.geode.cache.wan.GatewayReceiver
* @see org.apache.geode.cache.wan.GatewayEventFilter
* @see org.apache.geode.cache.wan.GatewayTransportFilter
* @see org.apache.geode.cache.wan.GatewaySender.OrderPolicy
* @see org.apache.geode.cache.wan.GatewayEventSubstitutionFilter
* @since 2.2.0
*/
@Target(ElementType.TYPE)
@Retention(RetentionPolicy.RUNTIME)
@Inherited
@Documented
@Import({ GatewaySenderBeanFactoryPostProcessor.class, GatewaySendersConfiguration.class })
@SuppressWarnings("unused")
public @interface EnableGatewaySenders {
/**
* This property configures the time, in milliseconds, that an object can be in the queue to be replicated before the
* {@link org.apache.geode.cache.wan.GatewaySender} logs an alert.
* <p> This property can also be configured using the <b><i>spring.data.gemfire.gateway.sender.alert-threshold</i></b>
* <p>Default value is {@link GatewaySenderConfiguration#DEFAULT_ALERT_THRESHOLD}
*/
int alertThreshold() default GatewaySenderConfiguration.DEFAULT_ALERT_THRESHOLD;
/**
* A boolean flag to indicate if the configured GatewaySender(s) should use conflate entries in each batch. This means,
* that a batch will never contain duplicate entries, as the batch will always only contain the latest value for a key.
* <p> This property can also be configured using the <b><i>spring.data.gemfire.gateway.sender.batch-conflation-enabled</i></b>
* <p>Default value is {@value GatewaySenderConfiguration#DEFAULT_BATCH_CONFLATION_ENABLED}
*/
boolean batchConflationEnabled() default GatewaySenderConfiguration.DEFAULT_BATCH_CONFLATION_ENABLED;
/**
* This property configures the maximum batch size that the {@link org.apache.geode.cache.wan.GatewaySender} send to the
* remote site. This property works in conjunction with the {@link EnableGatewaySenders#batchTimeInterval()} setting.
* The {@link org.apache.geode.cache.wan.GatewaySender} will send when either the <b><i>batchSize</i></b> or <b><i>batchTimeInterval</i></b>
* is met.
* <p> This property can also be configured using the <b><i>spring.data.gemfire.gateway.sender.batch-size</i></b>
* <p>Default value is {@value GatewaySenderConfiguration#DEFAULT_BATCH_SIZE}
*/
int batchSize() default GatewaySenderConfiguration.DEFAULT_BATCH_SIZE;
/**
* This property configures the maximum batch time interval, in milliseconds, that the {@link org.apache.geode.cache.wan.GatewaySender} wait before
* attempting to send a batch of queued object to the remote {@link org.apache.geode.cache.wan.GatewayReceiver}.
* <p>This property works in conjunction with the {@link EnableGatewaySenders#batchSize()} setting.
* The {@link org.apache.geode.cache.wan.GatewaySender} will send when either the <b><i>batchSize</i></b> or <b><i>batchTimeInterval</i></b>
* is met.
* <p> This property can also be configured using the <b><i>spring.data.gemfire.gateway.sender.batch-time-interval</i></b>
* <p>Default value is {@value GatewaySenderConfiguration#DEFAULT_BATCH_TIME_INTERVAL}
*/
int batchTimeInterval() default GatewaySenderConfiguration.DEFAULT_BATCH_TIME_INTERVAL;
/**
* This property configures what {@link org.apache.geode.cache.DiskStore} the GatewaySender(s) are to use when persisting
* the GatewaySender(s) queues. This setting should be set when the <b><i>persistent</i></b> property is set to <b><i>true</i></b>.
* <p> This property can also be configured using the <b><i>spring.data.gemfire.gateway.sender.diskstore-reference</i></b>
* <p>Default value is {@value GatewaySenderConfiguration#DEFAULT_DISK_STORE_REFERENCE}
*/
String diskStoreReference() default GatewaySenderConfiguration.DEFAULT_DISK_STORE_REFERENCE;
/**
* A boolean flag to indicate if the configured GatewaySender(s) should use synchronous diskstore writes.
* <p> This property can also be configured using the <b><i>spring.data.gemfire.gateway.sender.disk-synchronous</i></b>
* <p>Default value is {@value GatewaySenderConfiguration#DEFAULT_DISK_SYNCHRONOUS}
*/
boolean diskSynchronous() default GatewaySenderConfiguration.DEFAULT_DISK_SYNCHRONOUS;
/**
* This property configures the number of dispatcher threads that the {@link org.apache.geode.cache.wan.GatewaySender}
* will try to use to dispatch the queuing object.
* <p>This property can also be configured using the <b><i>spring.data.gemfire.gateway.sender.dispatcher-threads</i></b>
* <p>Default value is {@value GatewaySenderConfiguration#DEFAULT_DISPATCHER_THREADS}
*/
int dispatcherThreads() default GatewaySenderConfiguration.DEFAULT_DISPATCHER_THREADS;
/**
* A list of {@link org.apache.geode.cache.wan.GatewayEventFilter} to be applied to the {@link org.apache.geode.cache.wan.GatewaySender}.
* {@link org.apache.geode.cache.wan.GatewayEventFilter} are used to filter out objects from the sending queue before dispatching
* them to the remote {@link org.apache.geode.cache.wan.GatewayReceiver}.
* <p>This property can also be configured using the <b><i>spring.data.gemfire.gateway.sender.event-filters</i></b> property.
* <p>Default value is and empty list
*/
String[] eventFilters() default {};
/**
* This property configures the {@link org.apache.geode.cache.wan.GatewayEventSubstitutionFilter} to be used by the GatewaySender(s).
* The {@link org.apache.geode.cache.wan.GatewayEventSubstitutionFilter} is used to replace values on objects before they
* are enqueue for remote replication.
* <p> This property can also be configured using the <b><i>spring.data.gemfire.gateway.sender.event-substitution-filter</i></b>
* <p>Default value is {@link GatewaySenderConfiguration#DEFAULT_EVENT_SUBSTITUTION_FILTER}
*/
String eventSubstitutionFilter() default GatewaySenderConfiguration.DEFAULT_EVENT_SUBSTITUTION_FILTER;
/**
* A list of {@link org.apache.geode.cache.wan.GatewaySender} to be configured. If no GatewaySenders are configured,
* a default GatewaySender will be created using the properties provided.
*/
EnableGatewaySender[] gatewaySenders() default {};
/**
* A boolean to indicate if a configured {@link org.apache.geode.cache.wan.GatewaySender} should be started automatically
* <p>This property can also be configured using the <b><i>spring.data.gemfire.gateway.sender.manual-start</i></b>
* <p>Default is {@value @EnableGatewaySenderConfiguration.DEFAULT_MANUAL_START}
*/
boolean manualStart() default GatewaySenderConfiguration.DEFAULT_MANUAL_START;
/**
* This property configures the maximum size, in MB, that the {@link org.apache.geode.cache.wan.GatewaySender} that the queue
* may take on heap memory, before overflowing to disk.
* <p>This property can also be configured using the <b><i>spring.data.gemfire.gateway.sender.maximum-queue-memory</i></b>
* <p>Default value is {@value GatewaySenderConfiguration#DEFAULT_MAXIMUM_QUEUE_MEMORY}
*/
int maximumQueueMemory() default GatewaySenderConfiguration.DEFAULT_MAXIMUM_QUEUE_MEMORY;
/**
* The id of the distributed system (cluster) that the GatewaySender(s) would send their data to.
* <p>This property can also be configured using the <b><i>spring.data.gemfire.gateway.sender.remote-distributed-system-id</i></b>
* <p>Default is {@value @EnableGatewaySenderConfiguration.DEFAULT_REMOTE_DISTRIBUTED_SYSTEM_ID}
*/
int remoteDistributedSystemId() default GatewaySenderConfiguration.DEFAULT_REMOTE_DISTRIBUTED_SYSTEM_ID;
/**
* This property sets the ordering policy that the GatewaySender(s) will use when queueing entries to be replicated to
* a remote {@link org.apache.geode.cache.wan.GatewayReceiver}.
* <p>There are three different ordering policies:
* <ul>
* <li>{@link OrderPolicyType#KEY} - Order of events preserved on a per key basis</li>
* <li>{@link OrderPolicyType#THREAD} - Order of events preserved by the thread that added the event</li>
* <li>{@link OrderPolicyType#PARTITION} - Order of events is preserved in order that they arrived in partitioned Region</li>
* </ul>
* <p>This property can also be configured using the <b><i>spring.data.gemfire.gateway.sender.order-policy</i></b></p>
* <p>Default value is an empty list
*/
OrderPolicyType orderPolicy() default OrderPolicyType.KEY;
/**
* A boolean flag to indicate if the configured GatewaySender(s) should use parallel GatewaySender replication.
* Parallel replication means that each {@link org.apache.geode.cache.server.CacheServer} that defines a {@link org.apache.geode.cache.wan.GatewaySender}
* will send data to a remote {@link org.apache.geode.cache.wan.GatewayReceiver}.
* <p> This property can also be configured using the <b><i>spring.data.gemfire.gateway.sender.parallel</i></b>
* <p>Default value is {@value GatewaySenderConfiguration#DEFAULT_PARALLEL}
*/
boolean parallel() default GatewaySenderConfiguration.DEFAULT_PARALLEL;
/**
* A boolean flag to indicate if the configured GatewaySender(s) should use persistence. This setting should be used
* in conjunction with the <b><i>diskStoreReference</i></b> property.
* <p> This property can also be configured using the <b><i>spring.data.gemfire.gateway.sender.persistent</i></b>
* <p>Default value is {@value GatewaySenderConfiguration#DEFAULT_PERSISTENT}
*/
boolean persistent() default GatewaySenderConfiguration.DEFAULT_PERSISTENT;
/**
* A list of {@link org.apache.geode.cache.Region} names that are to be configured with {@link org.apache.geode.cache.wan.GatewaySender} replication.
* An empty list will denote that ALL regions are to be replicated to the remote {@link org.apache.geode.cache.wan.GatewayReceiver}.
* <p>This property can also be configured using the <b><i>spring.data.gemfire.gateway.sender.region-names</i></b> property.
* <p>Default value is an empty list
*/
String[] regions() default {};
/**
* The socket buffer size for the {@link org.apache.geode.cache.wan.GatewaySender}. This setting is in bytes.
* <p>This property can also be configured using the <b><i>spring.data.gemfire.gateway.sender.socket-buffer-size</i></b> property.
* <p>Default value is {@value GatewaySenderConfiguration#DEFAULT_SOCKET_BUFFER_SIZE}
*/
int socketBufferSize() default GatewaySenderConfiguration.DEFAULT_SOCKET_BUFFER_SIZE;
/**
* Amount of time in milliseconds that the gateway sender will wait to receive an acknowledgment from a remote site.
* By default this is set to 0, which means there is no timeout. The minimum allowed timeout is 30000 (milliseconds).
* <p>This property can also be configured using the <b><i>spring.data.gemfire.gateway.sender.socket-read-timeout</i></b> property.
* <p>Default value is {@value GatewaySenderConfiguration#DEFAULT_SOCKET_READ_TIMEOUT}
*/
int socketReadTimeout() default GatewaySenderConfiguration.DEFAULT_SOCKET_READ_TIMEOUT;
/**
* An in-order list of {@link org.apache.geode.cache.wan.GatewayTransportFilter} to be applied to
* {@link org.apache.geode.cache.wan.GatewaySender}
* <p>This property can also be configured using the <b><i>spring.data.gemfire.gateway.sender.transport-filters</i></b>s property
* <p>Default value is an empty list
*/
String[] transportFilters() default {};
}

View File

@@ -0,0 +1,517 @@
package org.springframework.data.gemfire.config.annotation;
import java.lang.annotation.Annotation;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.function.Supplier;
import org.apache.geode.cache.wan.GatewaySender;
import org.springframework.beans.MutablePropertyValues;
import org.springframework.beans.factory.config.BeanReference;
import org.springframework.beans.factory.config.RuntimeBeanReference;
import org.springframework.beans.factory.support.BeanDefinitionBuilder;
import org.springframework.beans.factory.support.BeanDefinitionRegistry;
import org.springframework.beans.factory.support.ManagedList;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.ImportAware;
import org.springframework.context.annotation.ImportBeanDefinitionRegistrar;
import org.springframework.core.annotation.AnnotationAttributes;
import org.springframework.core.type.AnnotationMetadata;
import org.springframework.data.gemfire.config.annotation.support.AbstractAnnotationConfigSupport;
import org.springframework.data.gemfire.wan.GatewaySenderFactoryBean;
import org.springframework.data.gemfire.wan.OrderPolicyType;
import org.springframework.util.StringUtils;
@Configuration
public class GatewaySenderConfiguration extends AbstractAnnotationConfigSupport
implements ImportBeanDefinitionRegistrar, ImportAware {
public static final String DEFAULT_NAME = "GatewaySender";
public static final int DEFAULT_SOCKET_BUFFER_SIZE = GatewaySender.DEFAULT_SOCKET_BUFFER_SIZE;
static final boolean DEFAULT_MANUAL_START = GatewaySender.DEFAULT_MANUAL_START;
static final int DEFAULT_REMOTE_DISTRIBUTED_SYSTEM_ID = GatewaySender.DEFAULT_DISTRIBUTED_SYSTEM_ID;
static final boolean DEFAULT_DISK_SYNCHRONOUS = GatewaySender.DEFAULT_DISK_SYNCHRONOUS;
static final boolean DEFAULT_BATCH_CONFLATION_ENABLED = GatewaySender.DEFAULT_BATCH_CONFLATION;
static final boolean DEFAULT_PARALLEL = GatewaySender.DEFAULT_IS_PARALLEL;
static final boolean DEFAULT_PERSISTENT = GatewaySender.DEFAULT_PERSISTENCE_ENABLED;
static final int DEFAULT_ALERT_THRESHOLD = GatewaySender.DEFAULT_ALERT_THRESHOLD;
static final int DEFAULT_BATCH_SIZE = GatewaySender.DEFAULT_BATCH_SIZE;
static final int DEFAULT_BATCH_TIME_INTERVAL = GatewaySender.DEFAULT_BATCH_TIME_INTERVAL;
static final int DEFAULT_DISPATCHER_THREADS = GatewaySender.DEFAULT_DISPATCHER_THREADS;
static final int DEFAULT_MAXIMUM_QUEUE_MEMORY = GatewaySender.DEFAULT_MAXIMUM_QUEUE_MEMORY;
static final int DEFAULT_SOCKET_READ_TIMEOUT = 0;
static final OrderPolicyType DEFAULT_ORDER_POLICY = OrderPolicyType.KEY;
static final String DEFAULT_EVENT_SUBSTITUTION_FILTER = "";
static final String[] DEFAULT_EVENT_FILTERS = {};
static final String[] DEFAULT_TRANSPORT_FILTERS = {};
static final String[] DEFAULT_REGION_NAMES = {};
static final String DEFAULT_DISK_STORE_REFERENCE = "";
static final String REGION_NAMES_LITERAL = "regions";
private static final String NAME_LITERAL = "name";
private static final String MANUAL_START_LITERAL = "manualStart";
private static final String REMOTE_DIST_SYSTEM_ID_LITERAL = "remoteDistributedSystemId";
private static final String DISK_SYNCHRONOUS_LITERAL = "diskSynchronous";
private static final String BATCH_CONFLATION_ENABLED_LITERAL = "batchConflationEnabled";
private static final String PARALLEL_LITERAL = "parallel";
private static final String PERSISTENT_LITERAL = "persistent";
private static final String ORDER_POLICY_LITERAL = "orderPolicy";
private static final String EVENT_SUBSTITUTION_FILTER_LITERAL = "eventSubstitutionFilter";
private static final String ALERT_THRESHOLD_LITERAL = "alertThreshold";
private static final String BATCH_SIZE_LITERAL = "batchSize";
private static final String BATCH_TIME_INTERVAL_LITERAL = "batchTimeInterval";
private static final String DISPATCHER_THREAD_LITERAL = "dispatcherThreads";
private static final String MAXIMUM_QUEUE_MEMORY_LITERAL = "maximumQueueMemory";
private static final String SOCKET_BUFFER_SIZE_LITERAL = "socketBufferSize";
private static final String SOCKET_READ_TIMEOUT_LITERAL = "socketReadTimeout";
private static final String DISK_STORE_REFERENCE_LITERAL = "diskStoreReference";
private static final String EVENT_FILTERS_LITERAL = "eventFilters";
private static final String TRANSPORT_FILTERS_LITERAL = "transportFilters";
private static final String MANUAL_START_PROPERTY_NAME = "manual-start";
private static final String REMOTE_DIST_SYSTEM_ID_PROPERTY_NAME = "remote-distributed-system-id";
private static final String DISK_SYNCHRONOUS_PROPERTY_NAME = "disk-synchronous";
private static final String BATCH_CONFLATION_ENABLED_PROPERTY_NAME = "batch-conflation-enabled";
private static final String PARALLEL_PROPERTY_NAME = "parallel";
private static final String PERSISTENT_PROPERTY_NAME = "persistent";
private static final String ORDER_POLICY_PROPERTY_NAME = "order-policy";
private static final String EVENT_SUBSTITUTION_FILTER_PROPERTY_NAME = "event-substitution-filter";
private static final String ALERT_THRESHOLD_PROPERTY_NAME = "alert-threshold";
private static final String BATCH_SIZE_PROPERTY_NAME = "batch-size";
private static final String BATCH_TIME_INTERVAL_PROPERTY_NAME = "batch-time-interval";
private static final String DISPATCHER_THREAD_PROPERTY_NAME = "dispatcher-threads";
private static final String MAXIMUM_QUEUE_MEMORY_PROPERTY_NAME = "maximum-queue-memory";
private static final String SOCKET_BUFFER_SIZE_PROPERTY_NAME = "socket-buffer-size";
private static final String SOCKET_READ_TIMEOUT_PROPERTY_NAME = "socket-read-timeout";
private static final String DISK_STORE_REFERENCE_PROPERTY_NAME = "disk-store-reference";
private static final String EVENT_FILTERS_PROPERTY_NAME = "event-filters";
private static final String TRANSPORT_FILTERS_PROPERTY_NAME = "transport-filters";
private static final String REGION_NAMES_PROPERTY_NAME = "regions";
private String gatewaySenderBeanName;
private List<GatewaySenderConfigurer> gatewaySenderConfigurers = Collections.emptyList();
/**
* Processes the {@link EnableGatewaySender} annotation and registers the configured BeanDefinition
*
* @param annotationMetadata
* @param beanDefinitionRegistry
*/
@Override
public void registerBeanDefinitions(AnnotationMetadata annotationMetadata,
BeanDefinitionRegistry beanDefinitionRegistry) {
if (annotationMetadata.hasAnnotation(EnableGatewaySender.class.getName())) {
Map<String, Object> annotationAttributesMap = annotationMetadata
.getAnnotationAttributes(EnableGatewaySender.class.getName());
AnnotationAttributes gatewaySenderAnnotation = AnnotationAttributes.fromMap(annotationAttributesMap);
registerGatewaySender(gatewaySenderAnnotation, beanDefinitionRegistry, null);
}
}
@Override protected Class<? extends Annotation> getAnnotationType() {
return EnableGatewaySender.class;
}
/**
* This method processes a defined {@link EnableGatewaySender} on the {@link EnableGatewaySenders} annotation.
* It will process properties defined in either annotation or <b><i>application.properties</i></b>.
*
* @param gatewaySenderAnnotation
* @param registry
* @param parentGatewaySenderAnnotation
*/
protected void registerGatewaySender(String gatewaySenderName, AnnotationAttributes gatewaySenderAnnotation,
BeanDefinitionRegistry registry, AnnotationAttributes parentGatewaySenderAnnotation) {
BeanDefinitionBuilder gatewaySenderBuilder = BeanDefinitionBuilder.genericBeanDefinition(
GatewaySenderFactoryBean.class);
configureGatewaySenderFromAnnotation(gatewaySenderName, gatewaySenderAnnotation, gatewaySenderBuilder,
parentGatewaySenderAnnotation);
configureGatewaySenderFromProperties(gatewaySenderName, gatewaySenderBuilder);
configureGatewaySenderArguments(gatewaySenderName, gatewaySenderAnnotation, gatewaySenderBuilder,
parentGatewaySenderAnnotation);
registry.registerBeanDefinition(gatewaySenderName, gatewaySenderBuilder.getBeanDefinition());
}
private void configureGatewaySenderArguments(String gatewaySenderName,
AnnotationAttributes gatewaySenderAnnotation, BeanDefinitionBuilder gatewaySenderBuilder,
AnnotationAttributes parentGatewaySenderAnnotation) {
String[] resolveValue = Optional
.ofNullable(getStringArrayFromAnnotation(gatewaySenderAnnotation, REGION_NAMES_LITERAL))
.filter(childValue -> !childValue.equals(DEFAULT_REGION_NAMES))
.orElse(
Optional.ofNullable(getStringArrayFromAnnotation(parentGatewaySenderAnnotation, REGION_NAMES_LITERAL))
.orElse(DEFAULT_REGION_NAMES));
String[] resolvedPropertyValue = Optional
.ofNullable(resolveValueFromProperty(gatewaySenderName, REGION_NAMES_PROPERTY_NAME, resolveValue))
.orElse(resolveValue);
setPropertyValueOrUseParent(gatewaySenderBuilder, REGION_NAMES_LITERAL, resolvedPropertyValue, null,
DEFAULT_REGION_NAMES);
}
/**
* This method processes a defined {@link EnableGatewaySender} on the {@link EnableGatewaySenders} annotation.
* It will process properties defined in either annotation or <b><i>application.properties</i></b>.
*
* @param gatewaySenderAnnotation
* @param registry
* @param parentGatewaySenderAnnotation
*/
protected void registerGatewaySender(AnnotationAttributes gatewaySenderAnnotation,
BeanDefinitionRegistry registry, AnnotationAttributes parentGatewaySenderAnnotation) {
String gatewaySenderName = getStringFromAnnotation(gatewaySenderAnnotation, NAME_LITERAL);
registerGatewaySender(gatewaySenderName, gatewaySenderAnnotation, registry, parentGatewaySenderAnnotation);
}
public void setGatewaySenderBeanName(String gatewaySenderBeanName) {
this.gatewaySenderBeanName = gatewaySenderBeanName;
}
/**
* In this method a {@link GatewaySender} is configured from properties defined on the {@link EnableGatewaySender}
* annotation
*
* @param gatewaySenderAnnotation
* @param gatewaySenderBuilder
* @param parentGatewaySenderAnnotation
*/
private void configureGatewaySenderFromAnnotation(String gatewaySenderName,
AnnotationAttributes gatewaySenderAnnotation,
BeanDefinitionBuilder gatewaySenderBuilder,
AnnotationAttributes parentGatewaySenderAnnotation) {
setGatewaySenderBeanName(Optional.ofNullable(gatewaySenderName).orElse(DEFAULT_NAME));
setPropertyValueIfNotDefault(gatewaySenderBuilder, NAME_LITERAL,
gatewaySenderName, DEFAULT_NAME);
setPropertyValueOrUseParent(gatewaySenderBuilder, MANUAL_START_LITERAL,
getBooleanFromAnnotation(gatewaySenderAnnotation, MANUAL_START_LITERAL),
getBooleanFromAnnotation(parentGatewaySenderAnnotation, MANUAL_START_LITERAL), DEFAULT_MANUAL_START);
setPropertyValueOrUseParent(gatewaySenderBuilder, REMOTE_DIST_SYSTEM_ID_LITERAL,
getNumberFromAnnotation(gatewaySenderAnnotation, REMOTE_DIST_SYSTEM_ID_LITERAL),
getNumberFromAnnotation(parentGatewaySenderAnnotation, REMOTE_DIST_SYSTEM_ID_LITERAL),
DEFAULT_REMOTE_DISTRIBUTED_SYSTEM_ID);
setPropertyValueOrUseParent(gatewaySenderBuilder, DISK_STORE_REFERENCE_LITERAL,
getStringFromAnnotation(gatewaySenderAnnotation, DISK_STORE_REFERENCE_LITERAL),
getStringFromAnnotation(parentGatewaySenderAnnotation, DISK_STORE_REFERENCE_LITERAL),
DEFAULT_DISK_STORE_REFERENCE);
setPropertyValueOrUseParent(gatewaySenderBuilder, DISK_SYNCHRONOUS_LITERAL,
getBooleanFromAnnotation(gatewaySenderAnnotation, DISK_SYNCHRONOUS_LITERAL),
getBooleanFromAnnotation(parentGatewaySenderAnnotation, DISK_SYNCHRONOUS_LITERAL),
DEFAULT_DISK_SYNCHRONOUS);
setPropertyValueOrUseParent(gatewaySenderBuilder, BATCH_CONFLATION_ENABLED_LITERAL,
getBooleanFromAnnotation(gatewaySenderAnnotation, BATCH_CONFLATION_ENABLED_LITERAL),
getBooleanFromAnnotation(parentGatewaySenderAnnotation, BATCH_CONFLATION_ENABLED_LITERAL),
DEFAULT_BATCH_CONFLATION_ENABLED);
setPropertyValueOrUseParent(gatewaySenderBuilder, BATCH_SIZE_LITERAL,
getNumberFromAnnotation(gatewaySenderAnnotation, BATCH_SIZE_LITERAL),
getNumberFromAnnotation(parentGatewaySenderAnnotation, BATCH_SIZE_LITERAL), DEFAULT_BATCH_SIZE);
setPropertyValueOrUseParent(gatewaySenderBuilder, BATCH_TIME_INTERVAL_LITERAL,
getNumberFromAnnotation(gatewaySenderAnnotation, BATCH_TIME_INTERVAL_LITERAL),
getNumberFromAnnotation(parentGatewaySenderAnnotation, BATCH_TIME_INTERVAL_LITERAL),
DEFAULT_BATCH_TIME_INTERVAL);
setPropertyValueOrUseParent(gatewaySenderBuilder, PARALLEL_LITERAL,
getBooleanFromAnnotation(gatewaySenderAnnotation, PARALLEL_LITERAL),
getBooleanFromAnnotation(parentGatewaySenderAnnotation, PARALLEL_LITERAL), DEFAULT_PARALLEL);
setPropertyValueOrUseParent(gatewaySenderBuilder, PERSISTENT_LITERAL,
getBooleanFromAnnotation(gatewaySenderAnnotation, PERSISTENT_LITERAL),
getBooleanFromAnnotation(parentGatewaySenderAnnotation, PERSISTENT_LITERAL), DEFAULT_PERSISTENT);
setPropertyValueOrUseParent(gatewaySenderBuilder, ORDER_POLICY_LITERAL,
getEnumFromAnnotation(gatewaySenderAnnotation, ORDER_POLICY_LITERAL, DEFAULT_ORDER_POLICY).toString(),
getEnumFromAnnotation(parentGatewaySenderAnnotation, ORDER_POLICY_LITERAL, DEFAULT_ORDER_POLICY).toString(),
DEFAULT_ORDER_POLICY.toString());
setPropertyValueOrUseParentAsBeanReferenceList(gatewaySenderBuilder, EVENT_FILTERS_LITERAL,
getStringArrayFromAnnotation(gatewaySenderAnnotation, EVENT_FILTERS_LITERAL),
getStringArrayFromAnnotation(parentGatewaySenderAnnotation, EVENT_FILTERS_LITERAL), DEFAULT_EVENT_FILTERS);
setPropertyValueIfNotDefaultAsBeanReference(gatewaySenderBuilder, EVENT_SUBSTITUTION_FILTER_LITERAL,
getStringFromAnnotation(gatewaySenderAnnotation, EVENT_SUBSTITUTION_FILTER_LITERAL),
getStringFromAnnotation(parentGatewaySenderAnnotation, EVENT_SUBSTITUTION_FILTER_LITERAL),
DEFAULT_EVENT_SUBSTITUTION_FILTER);
setPropertyValueOrUseParent(gatewaySenderBuilder, ALERT_THRESHOLD_LITERAL,
getNumberFromAnnotation(gatewaySenderAnnotation, ALERT_THRESHOLD_LITERAL),
getNumberFromAnnotation(parentGatewaySenderAnnotation, ALERT_THRESHOLD_LITERAL), DEFAULT_ALERT_THRESHOLD);
setPropertyValueOrUseParent(gatewaySenderBuilder, DISPATCHER_THREAD_LITERAL,
getNumberFromAnnotation(gatewaySenderAnnotation, DISPATCHER_THREAD_LITERAL),
getNumberFromAnnotation(parentGatewaySenderAnnotation, DISPATCHER_THREAD_LITERAL),
DEFAULT_DISPATCHER_THREADS);
setPropertyValueOrUseParent(gatewaySenderBuilder, MAXIMUM_QUEUE_MEMORY_LITERAL,
getNumberFromAnnotation(gatewaySenderAnnotation, MAXIMUM_QUEUE_MEMORY_LITERAL),
getNumberFromAnnotation(parentGatewaySenderAnnotation, MAXIMUM_QUEUE_MEMORY_LITERAL),
DEFAULT_MAXIMUM_QUEUE_MEMORY);
setPropertyValueOrUseParent(gatewaySenderBuilder, SOCKET_BUFFER_SIZE_LITERAL,
getNumberFromAnnotation(gatewaySenderAnnotation, SOCKET_BUFFER_SIZE_LITERAL),
getNumberFromAnnotation(parentGatewaySenderAnnotation, SOCKET_BUFFER_SIZE_LITERAL),
DEFAULT_SOCKET_BUFFER_SIZE);
setPropertyValueOrUseParent(gatewaySenderBuilder, SOCKET_READ_TIMEOUT_LITERAL,
getNumberFromAnnotation(gatewaySenderAnnotation, SOCKET_READ_TIMEOUT_LITERAL),
getNumberFromAnnotation(parentGatewaySenderAnnotation, SOCKET_READ_TIMEOUT_LITERAL),
DEFAULT_SOCKET_READ_TIMEOUT);
setPropertyValueOrUseParentAsBeanReferenceList(gatewaySenderBuilder, TRANSPORT_FILTERS_LITERAL,
getStringArrayFromAnnotation(gatewaySenderAnnotation, TRANSPORT_FILTERS_LITERAL),
getStringArrayFromAnnotation(parentGatewaySenderAnnotation, TRANSPORT_FILTERS_LITERAL),
DEFAULT_TRANSPORT_FILTERS);
}
private Enum<?> getEnumFromAnnotation(AnnotationAttributes gatewaySenderAnnotation, String annotationLabel,
Enum<?> defaultValue) {
Enum returnValue = defaultValue;
if (gatewaySenderAnnotation != null) {
returnValue = getValueFromAnnotation(gatewaySenderAnnotation, () ->
gatewaySenderAnnotation.getEnum(annotationLabel));
}
return returnValue;
}
private String getStringFromAnnotation(AnnotationAttributes gatewaySenderAnnotation, String annotationLabel) {
return getValueFromAnnotation(gatewaySenderAnnotation,
() -> gatewaySenderAnnotation.getString(annotationLabel));
}
private Integer getNumberFromAnnotation(AnnotationAttributes gatewaySenderAnnotation, String annotationLabel) {
return getValueFromAnnotation(gatewaySenderAnnotation,
() -> gatewaySenderAnnotation.getNumber(annotationLabel));
}
private String[] getStringArrayFromAnnotation(AnnotationAttributes gatewaySenderAnnotation,
String annotationLabel) {
return getValueFromAnnotation(gatewaySenderAnnotation,
() -> gatewaySenderAnnotation.getStringArray(annotationLabel));
}
private Boolean getBooleanFromAnnotation(AnnotationAttributes gatewaySenderAnnotation, String annotationLabel) {
return getValueFromAnnotation(gatewaySenderAnnotation,
() -> gatewaySenderAnnotation.getBoolean(annotationLabel));
}
private <T> T getValueFromAnnotation(AnnotationAttributes gatewaySenderAnnotation, Supplier<T> supplier) {
if (gatewaySenderAnnotation == null) {
return null;
}
else {
return supplier.get();
}
}
/**
* In this method a {@link GatewaySender} is configured from properties defined within an <b><i>application.properties</i></b>
* file.
* These properties are <i>"named"</i> properties and will follow the following pattern:
* <b><i>{@literal spring.data.gemfire.gateway.sender.<name>.manual-start}</i></b>
*
* @param gatewaySenderBeanBuilder
*/
private void configureGatewaySenderFromProperties(String gatewaySenderName,
BeanDefinitionBuilder gatewaySenderBeanBuilder) {
MutablePropertyValues beanPropertyValues = gatewaySenderBeanBuilder.getRawBeanDefinition().getPropertyValues();
gatewaySenderBeanBuilder.addPropertyValue("gatewaySenderConfigurers", resolveGatewaySenderConfigurers());
configureGatewaySenderFromProperty(gatewaySenderBeanBuilder, gatewaySenderName,
MANUAL_START_PROPERTY_NAME, MANUAL_START_LITERAL, DEFAULT_MANUAL_START);
configureGatewaySenderFromProperty(gatewaySenderBeanBuilder, gatewaySenderName,
REMOTE_DIST_SYSTEM_ID_PROPERTY_NAME, REMOTE_DIST_SYSTEM_ID_LITERAL, DEFAULT_REMOTE_DISTRIBUTED_SYSTEM_ID);
configureGatewaySenderFromProperty(gatewaySenderBeanBuilder, gatewaySenderName,
DISK_STORE_REFERENCE_PROPERTY_NAME, DISK_STORE_REFERENCE_LITERAL, new String[0]);
configureGatewaySenderFromProperty(gatewaySenderBeanBuilder, gatewaySenderName,
DISK_SYNCHRONOUS_PROPERTY_NAME, DISK_SYNCHRONOUS_LITERAL, DEFAULT_DISK_SYNCHRONOUS);
configureGatewaySenderFromProperty(gatewaySenderBeanBuilder, gatewaySenderName,
BATCH_CONFLATION_ENABLED_PROPERTY_NAME, BATCH_CONFLATION_ENABLED_LITERAL, DEFAULT_BATCH_CONFLATION_ENABLED);
configureGatewaySenderFromProperty(gatewaySenderBeanBuilder, gatewaySenderName,
BATCH_SIZE_PROPERTY_NAME, BATCH_SIZE_LITERAL, DEFAULT_BATCH_SIZE);
configureGatewaySenderFromProperty(gatewaySenderBeanBuilder, gatewaySenderName,
BATCH_TIME_INTERVAL_PROPERTY_NAME, BATCH_TIME_INTERVAL_LITERAL, DEFAULT_BATCH_TIME_INTERVAL);
configureGatewaySenderFromProperty(gatewaySenderBeanBuilder, gatewaySenderName,
PARALLEL_PROPERTY_NAME, PARALLEL_LITERAL, DEFAULT_PARALLEL);
configureGatewaySenderFromProperty(gatewaySenderBeanBuilder, gatewaySenderName,
PERSISTENT_PROPERTY_NAME, PERSISTENT_LITERAL, DEFAULT_PERSISTENT);
configureGatewaySenderFromProperty(gatewaySenderBeanBuilder, gatewaySenderName,
ORDER_POLICY_PROPERTY_NAME, ORDER_POLICY_LITERAL, DEFAULT_ORDER_POLICY.toString());
configureBeanReferenceListFromProperty(gatewaySenderBeanBuilder, gatewaySenderName,
EVENT_FILTERS_PROPERTY_NAME, EVENT_FILTERS_LITERAL, DEFAULT_EVENT_FILTERS);
configureBeanReferenceFromProperty(gatewaySenderBeanBuilder, gatewaySenderName,
EVENT_SUBSTITUTION_FILTER_PROPERTY_NAME, EVENT_SUBSTITUTION_FILTER_LITERAL,
DEFAULT_EVENT_SUBSTITUTION_FILTER);
configureGatewaySenderFromProperty(gatewaySenderBeanBuilder, gatewaySenderName,
ALERT_THRESHOLD_PROPERTY_NAME, ALERT_THRESHOLD_LITERAL, DEFAULT_ALERT_THRESHOLD);
configureGatewaySenderFromProperty(gatewaySenderBeanBuilder, gatewaySenderName,
DISPATCHER_THREAD_PROPERTY_NAME, DISPATCHER_THREAD_LITERAL, DEFAULT_DISPATCHER_THREADS);
configureGatewaySenderFromProperty(gatewaySenderBeanBuilder, gatewaySenderName,
MAXIMUM_QUEUE_MEMORY_PROPERTY_NAME, MAXIMUM_QUEUE_MEMORY_LITERAL, DEFAULT_MAXIMUM_QUEUE_MEMORY);
configureGatewaySenderFromProperty(gatewaySenderBeanBuilder, gatewaySenderName,
SOCKET_BUFFER_SIZE_PROPERTY_NAME, SOCKET_BUFFER_SIZE_LITERAL, DEFAULT_SOCKET_BUFFER_SIZE);
configureGatewaySenderFromProperty(gatewaySenderBeanBuilder, gatewaySenderName,
SOCKET_READ_TIMEOUT_PROPERTY_NAME, SOCKET_READ_TIMEOUT_LITERAL, DEFAULT_SOCKET_READ_TIMEOUT);
configureBeanReferenceListFromProperty(gatewaySenderBeanBuilder, gatewaySenderName,
TRANSPORT_FILTERS_PROPERTY_NAME, TRANSPORT_FILTERS_LITERAL, DEFAULT_TRANSPORT_FILTERS);
}
/**
* In this method a {@link GatewaySender} is configured from properties defined within an <b><i>application.properties</i></b>
* file.
* These properties are <i>"named"</i> properties and will follow the following pattern:
* <b><i>{@literal spring.data.gemfire.gateway.sender.<name>.manual-start}</i></b>
*
* @param gatewaySenderBeanBuilder
* @param gatewaySenderName
* @param propertyName
* @param defaultValue
*/
private <T extends Object> void configureGatewaySenderFromProperty(BeanDefinitionBuilder gatewaySenderBeanBuilder,
String gatewaySenderName, String propertyName, String fieldName, T defaultValue) {
T resolvedPropertyValue = resolveValueFromProperty(gatewaySenderName, propertyName, defaultValue);
if (resolvedPropertyValue == null) {
return;
}
setPropertyValueOrUseParent(gatewaySenderBeanBuilder, fieldName, resolvedPropertyValue, null, defaultValue);
}
private void configureBeanReferenceListFromProperty(BeanDefinitionBuilder gatewaySenderBeanBuilder,
String gatewaySenderName, String propertyName, String fieldName, String[] defaultValue) {
String[] resolvedPropertyValue = resolveValueFromProperty(gatewaySenderName, propertyName, defaultValue);
if (resolvedPropertyValue == null) {
return;
}
setPropertyValueOrUseParentAsBeanReferenceList(gatewaySenderBeanBuilder, fieldName, resolvedPropertyValue,
new String[] {},
defaultValue);
}
private void configureBeanReferenceFromProperty(BeanDefinitionBuilder gatewaySenderBeanBuilder,
String gatewaySenderName, String propertyName, String fieldName, String defaultValue) {
String resolvedPropertyValue = resolveValueFromProperty(gatewaySenderName, propertyName, defaultValue);
if (resolvedPropertyValue == null) {
return;
}
setPropertyValueIfNotDefaultAsBeanReference(gatewaySenderBeanBuilder, fieldName, resolvedPropertyValue, null,
defaultValue);
}
private List<GatewaySenderConfigurer> resolveGatewaySenderConfigurers() {
return Optional.ofNullable(this.gatewaySenderConfigurers)
.filter(gatewaySenderConfigurers -> !gatewaySenderConfigurers.isEmpty())
.orElseGet(() ->
Collections.singletonList(LazyResolvingComposableGatewaySenderConfigurer.create(getBeanFactory())));
}
private <T extends Object> T resolveValueFromProperty(String gatewaySenderName, String propertyName,
T defaultValue) {
Class<T> clazz = (Class<T>) defaultValue.getClass();
Optional<T> gatewaySenderProperty = Optional
.ofNullable(resolveProperty(gatewaySenderProperty(propertyName), clazz, null));
Optional<T> namedGatewaySenderProperty = Optional.ofNullable(
resolveProperty(namedGatewaySenderProperty(gatewaySenderName, propertyName), clazz, null));
return namedGatewaySenderProperty.orElse(gatewaySenderProperty.orElse(null));
}
private <T> BeanDefinitionBuilder setPropertyValueIfNotDefault(BeanDefinitionBuilder beanDefinitionBuilder,
String propertyName, T value, T defaultValue) {
return beanDefinitionBuilder.addPropertyValue(propertyName, Optional.ofNullable(value).orElse(defaultValue));
}
private <T> BeanDefinitionBuilder setPropertyValueOrUseParent(BeanDefinitionBuilder beanDefinitionBuilder,
String propertyName, T value, T parentValue, T defaultValue) {
T resolveValue = Optional.ofNullable(value).filter(childValue -> !childValue.equals(defaultValue))
.orElse(Optional.ofNullable(parentValue).orElse(defaultValue));
return beanDefinitionBuilder.addPropertyValue(propertyName, resolveValue);
}
private BeanDefinitionBuilder setPropertyValueIfNotDefaultAsBeanReference(
BeanDefinitionBuilder beanDefinitionBuilder,
String propertyName, String value, String parentValue, String defaultValue) {
if (!StringUtils.isEmpty(value)) {
return beanDefinitionBuilder.addPropertyReference(propertyName, value);
}
else if (!StringUtils.isEmpty(parentValue)) {
return beanDefinitionBuilder.addPropertyReference(propertyName, parentValue);
}
else {
if (!StringUtils.isEmpty(defaultValue)) {
beanDefinitionBuilder.addPropertyReference(propertyName, defaultValue);
}
}
return beanDefinitionBuilder;
}
private BeanDefinitionBuilder setPropertyValueOrUseParentAsBeanReferenceList(
BeanDefinitionBuilder beanDefinitionBuilder, String propertyName, String[] values, String[] parentValues,
String[] defaultList) {
ManagedList<BeanReference> beanReferences = new ManagedList<>();
String[] resolvedList = Optional.ofNullable(values).filter(t -> !Arrays.equals(t, defaultList))
.orElse(Optional.ofNullable(parentValues).orElse(defaultList));
Arrays.stream(resolvedList)
.map(RuntimeBeanReference::new)
.forEach(beanReferences::add);
return beanDefinitionBuilder.addPropertyValue(propertyName, beanReferences);
}
@Override public void setImportMetadata(AnnotationMetadata annotationMetadata) {
}
}

View File

@@ -0,0 +1,18 @@
package org.springframework.data.gemfire.config.annotation;
import org.springframework.data.gemfire.config.annotation.support.Configurer;
import org.springframework.data.gemfire.wan.GatewaySenderFactoryBean;
public interface GatewaySenderConfigurer extends Configurer<GatewaySenderFactoryBean> {
/**
* Configuration callback method providing a reference to a {@link org.springframework.data.gemfire.wan.GatewaySenderFactoryBean} used to construct,
* configure and initialize an instance of {@link org.apache.geode.cache.wan.GatewaySender}.
*
* @param beanName name of the {@link org.apache.geode.cache.wan.GatewaySender} bean declared in the Spring application context.
* @param bean reference to the {@link GatewaySenderFactoryBean}.
* @see org.springframework.data.gemfire.wan.GatewaySenderFactoryBean
* @see org.apache.geode.cache.wan.GatewaySender
*/
void configure(String beanName, GatewaySenderFactoryBean bean);
}

View File

@@ -0,0 +1,74 @@
package org.springframework.data.gemfire.config.annotation;
import java.lang.annotation.Annotation;
import java.util.Map;
import org.springframework.beans.factory.support.BeanDefinitionRegistry;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.annotation.AnnotationAttributes;
import org.springframework.core.type.AnnotationMetadata;
/**
* Spring {@link Configuration} class used to construct, configure and initialize {@link org.apache.geode.cache.wan.GatewaySender} instances
* in a Spring application context.
*
* @author Udo Kohlmeyer
* @author John Blum
* @see EnableGatewaySender
* @see EnableGatewaySenders
* @see org.apache.geode.cache.wan.GatewayReceiver
* @see org.springframework.context.annotation.Bean
* @see org.springframework.context.annotation.Configuration
* @since 2.2.0
*/
@Configuration
public class GatewaySendersConfiguration extends GatewaySenderConfiguration {
/**
* @param annotationMetadata
* @param beanDefinitionRegistry
*/
@Override
public void registerBeanDefinitions(AnnotationMetadata annotationMetadata,
BeanDefinitionRegistry beanDefinitionRegistry) {
if (annotationMetadata.hasAnnotation(EnableGatewaySenders.class.getName())) {
Map<String, Object> annotationAttributesMap = annotationMetadata
.getAnnotationAttributes(EnableGatewaySenders.class.getName());
AnnotationAttributes annotationAttributes = AnnotationAttributes.fromMap(annotationAttributesMap);
registerGatewaySenders(annotationAttributes, beanDefinitionRegistry);
}
}
/**
* @param gatewaySendersAnnotation
* @param registry
*/
private void registerGatewaySenders(AnnotationAttributes gatewaySendersAnnotation,
BeanDefinitionRegistry registry) {
AnnotationAttributes[] gatewaySenders = gatewaySendersAnnotation.getAnnotationArray("gatewaySenders");
if (gatewaySenders.length == 0) {
registerDefaultGatewaySender(registry, gatewaySendersAnnotation);
}
else {
for (AnnotationAttributes gatewaySender : gatewaySenders) {
registerGatewaySender(gatewaySender, registry, gatewaySendersAnnotation);
}
}
}
@Override protected Class<? extends Annotation> getAnnotationType() {
return EnableGatewaySenders.class;
}
private void registerDefaultGatewaySender(BeanDefinitionRegistry registry,
AnnotationAttributes parentGatewaySendersAnnotation) {
registerGatewaySender("GatewaySender", parentGatewaySendersAnnotation, registry,
parentGatewaySendersAnnotation);
}
}

View File

@@ -0,0 +1,49 @@
/*
* Copyright 2018 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
*
* https://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.data.gemfire.config.annotation;
import org.springframework.beans.factory.BeanFactory;
import org.springframework.data.gemfire.config.annotation.support.AbstractLazyResolvingComposableConfigurer;
import org.springframework.data.gemfire.wan.GatewaySenderFactoryBean;
import org.springframework.lang.Nullable;
/**
* Composition of {@link GatewaySenderConfigurer}.
*
* @author Udo Kohlmeyer
* @see GatewaySenderFactoryBean
* @see GatewaySenderConfigurer
* @see AbstractLazyResolvingComposableConfigurer
* @since 2.2.0
*/
public class LazyResolvingComposableGatewaySenderConfigurer
extends AbstractLazyResolvingComposableConfigurer<GatewaySenderFactoryBean, GatewaySenderConfigurer>
implements GatewaySenderConfigurer {
public static LazyResolvingComposableGatewaySenderConfigurer create() {
return create(null);
}
public static LazyResolvingComposableGatewaySenderConfigurer create(@Nullable BeanFactory beanFactory) {
return new LazyResolvingComposableGatewaySenderConfigurer().with(beanFactory);
}
@Override
protected Class<GatewaySenderConfigurer> getConfigurerType() {
return GatewaySenderConfigurer.class;
}
}

View File

@@ -778,6 +778,14 @@ public abstract class AbstractAnnotationConfigSupport
return String.format("%1$s%2$s", propertyName("gateway.receiver."), propertyNameSuffix);
}
protected String namedGatewaySenderProperty(String name, String propertyNameSuffix) {
return String.format("%1$s%2$s.%3$s", propertyName("gateway.sender."), name, propertyNameSuffix);
}
protected String gatewaySenderProperty(String propertyNameSuffix) {
return String.format("%1$s%2$s", propertyName("gateway.sender."), propertyNameSuffix);
}
/**
* Returns the fully-qualified {@link String property name}.
*

View File

@@ -0,0 +1,231 @@
package org.springframework.data.gemfire.config.support;
import static java.util.Arrays.stream;
import static org.springframework.data.gemfire.util.ArrayUtils.nullSafeArray;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collection;
import java.util.HashMap;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.stream.Collectors;
import org.springframework.beans.BeansException;
import org.springframework.beans.PropertyValue;
import org.springframework.beans.factory.annotation.AnnotatedBeanDefinition;
import org.springframework.beans.factory.config.BeanDefinition;
import org.springframework.beans.factory.config.BeanFactoryPostProcessor;
import org.springframework.beans.factory.config.ConfigurableListableBeanFactory;
import org.springframework.beans.factory.config.RuntimeBeanReference;
import org.springframework.beans.factory.support.ManagedList;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.type.MethodMetadata;
import org.springframework.data.gemfire.PeerRegionFactoryBean;
import org.springframework.data.gemfire.util.SpringUtils;
import org.springframework.data.gemfire.wan.GatewaySenderFactoryBean;
import org.springframework.util.StringUtils;
/**
* A {@link BeanFactoryPostProcessor} to associate the configured {@link org.apache.geode.cache.wan.GatewaySender}
* onto the corresponding {@link org.apache.geode.cache.Region}.
*
* @author Udo Kohlmeyer
* @since 2.2.0
*/
@Configuration
public class GatewaySenderBeanFactoryPostProcessor {
/**
* BeanFactory PostProcessor to map {@link GatewaySenderFactoryBean} to {@link org.apache.geode.cache.Region}
*
* @return BeanFactoryPostProcessor for the GatewaySenderFactoryBean
* @throws BeansException
*/
@Bean
public BeanFactoryPostProcessor postProcessBeanFactory() throws BeansException {
return beanFactory -> {
//Create a map of cached BeanDefinitions. Mapped under 'regions' and 'gatewaySenders' for easier lookups
Map<String, Map<String, BeanDefinition>> cachedBeanDefinitions = populateBeanDefinitionCache(beanFactory);
//Create a list of gatewaySender to Regions mapping
Map<String, List<String>> gatewaySenderToRegions = groupGatewaySenderPerRegion(cachedBeanDefinitions,
new HashMap<>());
//Add
addGatewaySendersToRegionFactory(beanFactory, gatewaySenderToRegions);
};
}
/**
* Determines if a {@link BeanDefinition} is of type {@link PeerRegionFactoryBean}, which means it is
* effectively of type {@link org.apache.geode.cache.Region}
* @param beanDefinition
* @param beanFactory
* @return {@literal boolean}
*/
private boolean isRegionBean(BeanDefinition beanDefinition, ConfigurableListableBeanFactory beanFactory) {
return Optional.ofNullable(beanDefinition)
.map(it -> resolveBeanClass(beanDefinition, beanFactory))
.filter(beanClass -> beanClass.isPresent())
.filter(beanClass -> PeerRegionFactoryBean.class.isAssignableFrom(beanClass.get()))
.isPresent();
}
private void addGatewaySendersToRegionFactory(ConfigurableListableBeanFactory beanFactory,
Map<String, List<String>> gatewaySenderToRegions) {
gatewaySenderToRegions.entrySet().forEach(entry -> {
List<RuntimeBeanReference> beanReferenceList = entry.getValue().stream()
.map(RuntimeBeanReference::new).collect(Collectors.toList());
ManagedList<RuntimeBeanReference> runtimeBeanReferences = new ManagedList<>();
runtimeBeanReferences.addAll(beanReferenceList);
Optional.ofNullable(beanFactory.getBeanDefinition(entry.getKey())).ifPresent(regionBeanDefinition ->
regionBeanDefinition.getPropertyValues().addPropertyValue("gatewaySenders", runtimeBeanReferences));
});
}
/**
* Mapping of GatewaySender to Regions.
* @param cachedBeanDefinitions
* @param gatewaySendersPerRegion
* @return
*/
private Map<String, List<String>> groupGatewaySenderPerRegion(
Map<String, Map<String, BeanDefinition>> cachedBeanDefinitions,
Map<String, List<String>> gatewaySendersPerRegion) {
Optional<Map<String, BeanDefinition>> gatewaySenders = Optional
.ofNullable(cachedBeanDefinitions.get("gatewaySenders"));
gatewaySenders.ifPresent(gatewaySendersMap -> gatewaySendersMap.forEach((key, gatewaySenderEntryValue) -> {
PropertyValue regions = gatewaySenderEntryValue.getPropertyValues().getPropertyValue("regions");
Optional<PropertyValue> regionsOptional = Optional.ofNullable(regions);
Collection<String> regionNames;
List<String> namedRegions = new ArrayList<>();
regionsOptional.ifPresent(valueHolder -> namedRegions.addAll(Arrays.asList((String[]) valueHolder.getValue())));
if (namedRegions.size() == 0) {
//Add gatewaySender to all Regions
regionNames = cachedBeanDefinitions.get("regions").keySet();
}
else {
//Add gatewaySender to named Regions
regionNames = namedRegions;
}
addGatewaySendersToRegion(gatewaySendersPerRegion, key,
regionNames);
}));
return gatewaySendersPerRegion;
}
/**
* Caches BeanDefinitions for GatewaySenders and Regions. Maps the beandefinitions under the keys of
* `regions` and `gatewaySenders`
* @param beanFactory
* @return Map containing all BeanDefinitions for Regions and GatewaySenders
*/
private Map<String, Map<String, BeanDefinition>> populateBeanDefinitionCache(
ConfigurableListableBeanFactory beanFactory) {
Map<String, Map<String, BeanDefinition>> cachedBeanDefinitions = new LinkedHashMap<>();
stream(nullSafeArray(beanFactory.getBeanDefinitionNames(), String.class))
.forEach(beanName -> Optional.of(beanFactory.getBeanDefinition(beanName))
.ifPresent(beanDefinition -> {
if (isRegionBean(beanDefinition, beanFactory)) {
addBeanDefinitionToList(beanName, beanDefinition, cachedBeanDefinitions, "regions");
}
else if (isGatewaySenderFactoryBean(beanDefinition)) {
addBeanDefinitionToList(beanName, beanDefinition, cachedBeanDefinitions, "gatewaySenders");
}
}));
return cachedBeanDefinitions;
}
/**
* Creates a map of <<String>,List<BeanDefinitions>> grouped by the `key`
*
* @param beanName
* @param beanDefinition
* @param cachedBeanDefinitions
* @param key
* @return
*/
private Map<String, Map<String, BeanDefinition>> addBeanDefinitionToList(String beanName,
BeanDefinition beanDefinition,
Map<String, Map<String, BeanDefinition>> cachedBeanDefinitions, String key) {
Map<String, BeanDefinition> beanDefinitions = cachedBeanDefinitions.get(key);
if (beanDefinitions == null) {
beanDefinitions = new HashMap<>();
}
beanDefinitions.put(beanName, beanDefinition);
cachedBeanDefinitions.put(key, beanDefinitions);
return cachedBeanDefinitions;
}
/**
* Mapping of gatewaySenders to individual regions. M -> N mapping capability
* @param regionMapping A map that holds the Region -> GatewaySender mapping
* @param gatewaySenderName The GatewaySender that is to be mapped to the regions
* @param regionNames Collection of defined regions that need to have gatewaySenders mapped to them.
*/
private void addGatewaySendersToRegion(Map<String, List<String>> regionMapping,
String gatewaySenderName, Collection<String> regionNames) {
regionNames.forEach(regionName -> {
List<String> gatewaySenders = regionMapping.get(regionName);
if (gatewaySenders == null) {
gatewaySenders = new ArrayList<>();
}
gatewaySenders.add(gatewaySenderName);
regionMapping.put(regionName, gatewaySenders);
});
}
private boolean isGatewaySenderFactoryBean(BeanDefinition beanDefinition) {
return GatewaySenderFactoryBean.class.getName()
.equals(beanDefinition.getBeanClassName());
}
protected Optional<Object> getPropertyValue(BeanDefinition beanDefinition, String propertyName) {
return SpringUtils.getPropertyValue(beanDefinition, propertyName);
}
/**
* Resolves the class type name of the bean defined by the given {@link BeanDefinition}.
*
* @param beanDefinition {@link BeanDefinition} defining the bean from which to resolve the class type name.
* @return an {@link Optional} {@link String} containing the resolved class type name of the bean defined
* by the given {@link BeanDefinition}.
* @see org.springframework.beans.factory.config.BeanDefinition#getBeanClassName()
*/
protected Optional<Class> resolveBeanClass(BeanDefinition beanDefinition,
ConfigurableListableBeanFactory beanFactory) {
Optional<String> beanClassName = Optional.ofNullable(beanDefinition)
.map(BeanDefinition::getBeanClassName)
.filter(StringUtils::hasText);
if (!beanClassName.isPresent()) {
beanClassName = Optional.ofNullable(beanDefinition)
.filter(it -> it instanceof AnnotatedBeanDefinition)
.filter(it -> StringUtils.hasText(it.getFactoryMethodName()))
.map(it -> ((AnnotatedBeanDefinition) it).getFactoryMethodMetadata())
.map(MethodMetadata::getReturnTypeName);
}
return beanClassName.map(className -> {
try {
return beanFactory.getBeanClassLoader().loadClass(className);
}
catch (ClassNotFoundException e) {
e.printStackTrace();
}
return null;
});
}
}

View File

@@ -16,16 +16,16 @@
package org.springframework.data.gemfire.wan;
import org.apache.geode.cache.Cache;
import org.apache.geode.cache.GemFireCache;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.DisposableBean;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.gemfire.support.AbstractFactoryBeanSupport;
import org.springframework.util.Assert;
import org.springframework.util.StringUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
/**
* Abstract base class for WAN Gateway components.
*
@@ -40,6 +40,7 @@ import org.slf4j.LoggerFactory;
public abstract class AbstractWANComponentFactoryBean<T> extends AbstractFactoryBeanSupport<T>
implements DisposableBean, InitializingBean {
@Autowired
protected Cache cache;
protected final Logger logger = LoggerFactory.getLogger(getClass());
@@ -49,12 +50,11 @@ public abstract class AbstractWANComponentFactoryBean<T> extends AbstractFactory
private String beanName;
private String name;
protected AbstractWANComponentFactoryBean(Cache cache) {
this.cache = cache;
protected AbstractWANComponentFactoryBean() {
}
public void setCache(Cache cache) {
this.cache = cache;
protected AbstractWANComponentFactoryBean(GemFireCache cache) {
this.cache = (Cache) cache;
}
@Override
@@ -70,6 +70,10 @@ public abstract class AbstractWANComponentFactoryBean<T> extends AbstractFactory
this.name = name;
}
public void setCache(Cache cache) {
this.cache = cache;
}
public String getName() {
return StringUtils.hasText(this.name)

View File

@@ -13,20 +13,29 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.data.gemfire.wan;
import static java.util.stream.StreamSupport.stream;
import static org.springframework.data.gemfire.util.CollectionUtils.nullSafeIterable;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
import java.util.Optional;
import org.apache.geode.cache.Cache;
import org.apache.geode.cache.GemFireCache;
import org.apache.geode.cache.wan.GatewayEventFilter;
import org.apache.geode.cache.wan.GatewayEventSubstitutionFilter;
import org.apache.geode.cache.wan.GatewaySender;
import org.apache.geode.cache.wan.GatewaySenderFactory;
import org.apache.geode.cache.wan.GatewayTransportFilter;
import org.apache.shiro.util.StringUtils;
import org.springframework.beans.factory.FactoryBean;
import org.springframework.data.gemfire.config.annotation.GatewaySenderConfigurer;
import org.springframework.data.gemfire.util.CollectionUtils;
/**
@@ -34,6 +43,7 @@ import org.springframework.data.gemfire.util.CollectionUtils;
*
* @author David Turanski
* @author John Blum
* @author Udo Kohlmeyer
* @see org.springframework.data.gemfire.wan.AbstractWANComponentFactoryBean
* @see org.apache.geode.cache.Cache
* @see org.apache.geode.cache.util.Gateway
@@ -73,6 +83,10 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
private String diskStoreReference;
private List<GatewaySenderConfigurer> gatewaySenderConfigurers = Collections.emptyList();
private List<String> regions;
/**
* Constructs an instance of the {@link GatewaySenderFactoryBean} class initialized with a reference to
* the Pivotal GemFire {@link Cache} used to configured and initialized a Pivotal GemFire {@link GatewaySender}.
@@ -80,8 +94,12 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
* @param cache reference to the Pivotal GemFire {@link Cache} used to create the Pivotal GemFire {@link GatewaySender}.
* @see org.apache.geode.cache.Cache
*/
public GatewaySenderFactoryBean(Cache cache) {
public GatewaySenderFactoryBean(GemFireCache cache) {
super(cache);
this.regions = Arrays.asList(new String[] {});
}
public GatewaySenderFactoryBean() {
}
/**
@@ -92,6 +110,9 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
GatewaySenderFactory gatewaySenderFactory = resolveGatewaySenderFactory();
stream(nullSafeIterable(gatewaySenderConfigurers).spliterator(), false)
.forEach(gatewaySenderConfigurer1 -> gatewaySenderConfigurer1.configure(getName(), this));
Optional.ofNullable(this.alertThreshold).ifPresent(gatewaySenderFactory::setAlertThreshold);
Optional.ofNullable(this.batchConflationEnabled).ifPresent(gatewaySenderFactory::setBatchConflationEnabled);
Optional.ofNullable(this.batchSize).ifPresent(gatewaySenderFactory::setBatchSize);
@@ -149,64 +170,60 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
: GatewaySender.class;
}
public void setAlertThreshold(Integer alertThreshold) {
this.alertThreshold = alertThreshold;
public boolean isManualStart() {
return manualStart;
}
public void setBatchConflationEnabled(Boolean batchConflationEnabled) {
this.batchConflationEnabled = batchConflationEnabled;
public void setManualStart(boolean manualStart) {
this.manualStart = manualStart;
}
/**
* Boolean value that determines whether Pivotal GemFire should conflate messages.
*
* @param enableBatchConflation a boolean value indicating whether Pivotal GemFire should conflate messages in the Queue.
* @see #setBatchConflationEnabled(Boolean)
* @deprecated use setBatchConflationEnabled(Boolean)
*/
@Deprecated
public void setEnableBatchConflation(Boolean enableBatchConflation) {
this.batchConflationEnabled = enableBatchConflation;
public int getRemoteDistributedSystemId() {
return remoteDistributedSystemId;
}
public void setBatchSize(Integer batchSize) {
this.batchSize = batchSize;
public void setRemoteDistributedSystemId(int remoteDistributedSystemId) {
this.remoteDistributedSystemId = remoteDistributedSystemId;
}
public void setBatchTimeInterval(Integer batchTimeInterval) {
this.batchTimeInterval = batchTimeInterval;
public GatewaySender getGatewaySender() {
return gatewaySender;
}
public void setDiskStoreRef(String diskStoreRef) {
this.diskStoreReference = diskStoreRef;
public void setGatewaySender(GatewaySender gatewaySender) {
this.gatewaySender = gatewaySender;
}
public List<GatewayEventFilter> getEventFilters() {
return eventFilters;
}
public void setEventFilters(List<GatewayEventFilter> eventFilters) {
this.eventFilters = eventFilters;
}
public List<GatewayTransportFilter> getTransportFilters() {
return transportFilters;
}
public void setTransportFilters(List<GatewayTransportFilter> transportFilters) {
this.transportFilters = transportFilters;
}
public Boolean getDiskSynchronous() {
return diskSynchronous;
}
public void setDiskSynchronous(Boolean diskSynchronous) {
this.diskSynchronous = diskSynchronous;
}
public void setDispatcherThreads(Integer dispatcherThreads) {
this.dispatcherThreads = dispatcherThreads;
}
public void setEventFilters(List<GatewayEventFilter> gatewayEventFilters) {
this.eventFilters = gatewayEventFilters;
}
public void setEventSubstitutionFilter(GatewayEventSubstitutionFilter eventSubstitutionFilter) {
this.eventSubstitutionFilter = eventSubstitutionFilter;
}
public void setManualStart(Boolean manualStart) {
this.manualStart = Boolean.TRUE.equals(manualStart);
}
public void setMaximumQueueMemory(Integer maximumQueueMemory) {
this.maximumQueueMemory = maximumQueueMemory;
}
public void setOrderPolicy(GatewaySender.OrderPolicy orderPolicy) {
this.orderPolicy = orderPolicy;
public void setBatchConflationEnabled(Boolean batchConflationEnabled) {
this.batchConflationEnabled = batchConflationEnabled;
}
public void setParallel(Boolean parallel) {
@@ -221,10 +238,6 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
return Boolean.TRUE.equals(this.parallel);
}
public void setPersistent(Boolean persistent) {
this.persistent = persistent;
}
public boolean isNotPersistent() {
return !isPersistent();
}
@@ -233,19 +246,107 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
return Boolean.TRUE.equals(this.persistent);
}
public void setRemoteDistributedSystemId(int remoteDistributedSystemId) {
this.remoteDistributedSystemId = remoteDistributedSystemId;
public void setPersistent(Boolean persistent) {
this.persistent = persistent;
}
public GatewaySender.OrderPolicy getOrderPolicy() {
return orderPolicy;
}
public void setOrderPolicy(GatewaySender.OrderPolicy orderPolicy) {
this.orderPolicy = orderPolicy;
}
public GatewayEventSubstitutionFilter getEventSubstitutionFilter() {
return eventSubstitutionFilter;
}
public void setEventSubstitutionFilter(GatewayEventSubstitutionFilter eventSubstitutionFilter) {
this.eventSubstitutionFilter = eventSubstitutionFilter;
}
public Integer getAlertThreshold() {
return alertThreshold;
}
public void setAlertThreshold(Integer alertThreshold) {
this.alertThreshold = alertThreshold;
}
public Integer getBatchSize() {
return batchSize;
}
public void setBatchSize(Integer batchSize) {
this.batchSize = batchSize;
}
public Integer getBatchTimeInterval() {
return batchTimeInterval;
}
public void setBatchTimeInterval(Integer batchTimeInterval) {
this.batchTimeInterval = batchTimeInterval;
}
public Integer getDispatcherThreads() {
return dispatcherThreads;
}
public void setDispatcherThreads(Integer dispatcherThreads) {
this.dispatcherThreads = dispatcherThreads;
}
public Integer getMaximumQueueMemory() {
return maximumQueueMemory;
}
public void setMaximumQueueMemory(Integer maximumQueueMemory) {
this.maximumQueueMemory = maximumQueueMemory;
}
public Integer getSocketBufferSize() {
return socketBufferSize;
}
public void setSocketBufferSize(Integer socketBufferSize) {
this.socketBufferSize = socketBufferSize;
}
public Integer getSocketReadTimeout() {
return socketReadTimeout;
}
public void setSocketReadTimeout(Integer socketReadTimeout) {
this.socketReadTimeout = socketReadTimeout;
}
public void setTransportFilters(List<GatewayTransportFilter> gatewayTransportFilters) {
this.transportFilters = gatewayTransportFilters;
public String getDiskStoreReference() {
return diskStoreReference;
}
public void setDiskStoreReference(String diskStoreReference) {
this.diskStoreReference = diskStoreReference;
}
public void setGatewaySenderConfigurers(List<GatewaySenderConfigurer> gatewaySenderConfigurers) {
this.gatewaySenderConfigurers = gatewaySenderConfigurers;
}
public void setDiskStoreRef(String diskStoreRef) {
this.diskStoreReference = diskStoreRef;
}
private List<String> getRegions() {
return regions;
}
public void setRegions(List<String> regions) {
this.regions = regions;
}
public void setRegions(String[] regions) {
this.regions = Arrays.asList(regions);
}
}

View File

@@ -13,10 +13,10 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.data.gemfire.wan;
import org.apache.geode.cache.util.Gateway;
import org.apache.geode.cache.wan.GatewaySender;
import org.springframework.data.gemfire.support.AbstractPropertyEditorConverterSupport;
/**
@@ -30,7 +30,7 @@ import org.springframework.data.gemfire.support.AbstractPropertyEditorConverterS
* @since 1.7.0
*/
@SuppressWarnings({ "deprecation", "unused" })
public class OrderPolicyConverter extends AbstractPropertyEditorConverterSupport<Gateway.OrderPolicy> {
public class OrderPolicyConverter extends AbstractPropertyEditorConverterSupport<GatewaySender.OrderPolicy> {
/**
* Converts the given String into a Pivotal GemFire Gateway.OrderPolicy enum.
@@ -43,9 +43,8 @@ public class OrderPolicyConverter extends AbstractPropertyEditorConverterSupport
* @see org.apache.geode.cache.util.Gateway.OrderPolicy
*/
@Override
public Gateway.OrderPolicy convert(final String source) {
public GatewaySender.OrderPolicy convert(String source) {
return assertConverted(source, OrderPolicyType.getOrderPolicy(OrderPolicyType.valueOfIgnoreCase(source)),
Gateway.OrderPolicy.class);
}
GatewaySender.OrderPolicy.class);
}
}

View File

@@ -13,10 +13,9 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.data.gemfire.wan;
import org.apache.geode.cache.util.Gateway;
import org.apache.geode.cache.wan.GatewaySender;
/**
* The OrderPolicyType class is an enumeration of Pivotal GemFire Gateway Order Policies.
@@ -28,11 +27,11 @@ import org.apache.geode.cache.util.Gateway;
@SuppressWarnings({ "deprecation", "unused" })
public enum OrderPolicyType {
KEY(Gateway.OrderPolicy.KEY),
PARTITION(Gateway.OrderPolicy.PARTITION),
THREAD(Gateway.OrderPolicy.THREAD);
KEY(GatewaySender.OrderPolicy.KEY),
PARTITION(GatewaySender.OrderPolicy.PARTITION),
THREAD(GatewaySender.OrderPolicy.THREAD);
private final Gateway.OrderPolicy orderPolicy;
private final GatewaySender.OrderPolicy orderPolicy;
/**
* Constructs an instance of the OrderPolicyType enum initialized with the matching Pivotal GemFire Gateway.OrderPolicy
@@ -41,7 +40,7 @@ public enum OrderPolicyType {
* @param orderPolicy the matching Pivotal GemFire Gateway.OrderPolicy enumerated value.
* @see org.apache.geode.cache.util.Gateway.OrderPolicy
*/
OrderPolicyType(final Gateway.OrderPolicy orderPolicy) {
OrderPolicyType(GatewaySender.OrderPolicy orderPolicy) {
this.orderPolicy = orderPolicy;
}
@@ -55,8 +54,8 @@ public enum OrderPolicyType {
* @see org.apache.geode.cache.util.Gateway.OrderPolicy
* @see #getOrderPolicy()
*/
public static Gateway.OrderPolicy getOrderPolicy(final OrderPolicyType orderPolicyType) {
return (orderPolicyType != null ? orderPolicyType.getOrderPolicy() : null);
public static GatewaySender.OrderPolicy getOrderPolicy(final OrderPolicyType orderPolicyType) {
return orderPolicyType != null ? orderPolicyType.getOrderPolicy() : null;
}
/**
@@ -68,7 +67,8 @@ public enum OrderPolicyType {
* @see org.apache.geode.cache.util.Gateway.OrderPolicy
* @see #getOrderPolicy()
*/
public static OrderPolicyType valueOf(final Gateway.OrderPolicy orderPolicy) {
public static OrderPolicyType valueOf(GatewaySender.OrderPolicy orderPolicy) {
for (OrderPolicyType orderPolicyType : values()) {
if (orderPolicyType.getOrderPolicy().equals(orderPolicy)) {
return orderPolicyType;
@@ -87,6 +87,7 @@ public enum OrderPolicyType {
* @see #name()
*/
public static OrderPolicyType valueOfIgnoreCase(final String name) {
for (OrderPolicyType orderPolicy : values()) {
if (orderPolicy.name().equalsIgnoreCase(name)) {
return orderPolicy;
@@ -102,7 +103,7 @@ public enum OrderPolicyType {
* @return a Pivotal GemFire Gateway.OrderPolicy for this enum.
* @see org.apache.geode.cache.util.Gateway.OrderPolicy
*/
public Gateway.OrderPolicy getOrderPolicy() {
public GatewaySender.OrderPolicy getOrderPolicy() {
return orderPolicy;
}
}

View File

@@ -0,0 +1,407 @@
package org.springframework.data.gemfire.config.annotation;
import static org.assertj.core.api.Assertions.assertThat;
import java.io.InputStream;
import java.io.OutputStream;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.TreeMap;
import java.util.stream.Collectors;
import org.apache.geode.cache.DataPolicy;
import org.apache.geode.cache.EntryEvent;
import org.apache.geode.cache.GemFireCache;
import org.apache.geode.cache.Region;
import org.apache.geode.cache.wan.GatewayEventFilter;
import org.apache.geode.cache.wan.GatewayEventSubstitutionFilter;
import org.apache.geode.cache.wan.GatewayQueueEvent;
import org.apache.geode.cache.wan.GatewaySender;
import org.apache.geode.cache.wan.GatewayTransportFilter;
import org.junit.After;
import org.junit.Before;
import org.junit.Test;
import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.context.annotation.AnnotationConfigApplicationContext;
import org.springframework.context.annotation.Bean;
import org.springframework.data.gemfire.PartitionedRegionFactoryBean;
import org.springframework.data.gemfire.test.mock.annotation.EnableGemFireMockObjects;
import org.springframework.data.gemfire.wan.GatewaySenderFactoryBean;
import org.springframework.data.gemfire.wan.OrderPolicyType;
/**
* Tests for {@link EnableGatewaySenders} and {@link EnableGatewaySender}.
*
* @author Udo Kohlmeyer
* @see org.junit.Test
* @see org.mockito.Mockito
* @see org.apache.geode.cache.server.CacheServer
* @see org.springframework.context.annotation.Configuration
* @see org.springframework.test.context.ContextConfiguration
* @see org.springframework.test.context.junit4.SpringRunner
* @see org.springframework.data.gemfire.wan.GatewayReceiverFactoryBean
* @see GatewaySenderConfigurer
* @see GatewaySenderConfiguration
* @since 2.2.0
*/
public class GatewaySenderConfigurationTests {
private ConfigurableApplicationContext applicationContext;
@Before
public void setup() {
}
@After
public void shutdown() {
Optional.ofNullable(this.applicationContext).ifPresent(ConfigurableApplicationContext::close);
}
@Test
public void annotationConfigurationOfMultipleGatewaySendersWithDefaultsFromParent() {
this.applicationContext = newApplicationContext(BaseGatewaySenderTestConfiguration.class,
TestConfigurationOfMultipleGatewaySenderAnnotationsButWithDefaultsFromParent.class);
TestGatewaySenderConfigurer gatewaySenderConfigurer = this.applicationContext
.getBean(TestGatewaySenderConfigurer.class);
Map<String, GatewaySender> beansOfType = this.applicationContext.getBeansOfType(GatewaySender.class);
assertThat(beansOfType.size()).isEqualTo(2);
String[] senders = new String[] { "TestGatewaySender", "TestGatewaySender2" };
assertThat(beansOfType.keySet().toArray()).containsExactly(senders);
for (String sender : senders) {
GatewaySender gatewaySender = this.applicationContext.getBean(sender, GatewaySender.class);
assertThat(gatewaySender.isManualStart()).isEqualTo(true);
assertThat(gatewaySender.getRemoteDSId()).isEqualTo(2);
assertThat(gatewaySender.getId()).isEqualTo(sender);
assertThat(gatewaySender.getDispatcherThreads()).isEqualTo(22);
assertThat(gatewaySender.isBatchConflationEnabled()).isEqualTo(true);
assertThat(gatewaySender.isParallel()).isEqualTo(true);
assertThat(gatewaySender.isPersistenceEnabled()).isEqualTo(true);
assertThat(gatewaySender.getDiskStoreName()).isEqualTo("someDiskStore");
assertThat(gatewaySender.getOrderPolicy()).isEqualTo(GatewaySender.OrderPolicy.PARTITION);
assertThat(gatewaySender.getGatewayEventSubstitutionFilter()).isNull();
assertThat(gatewaySender.getAlertThreshold()).isEqualTo(1234);
assertThat(gatewaySender.getBatchSize()).isEqualTo(1002);
assertThat(gatewaySender.getBatchTimeInterval()).isEqualTo(2000);
assertThat(gatewaySender.getMaximumQueueMemory()).isEqualTo(400);
assertThat(gatewaySender.getSocketReadTimeout()).isEqualTo(4000);
assertThat(gatewaySender.getSocketBufferSize()).isEqualTo(16384);
assertThat(gatewaySender.getGatewayTransportFilters().size()).isEqualTo(2);
assertThat(((GatewaySenderConfigurationTests.TestGatewayTransportFilter) gatewaySender
.getGatewayTransportFilters().get(0)).name).isEqualTo("transportBean2");
assertThat(((GatewaySenderConfigurationTests.TestGatewayTransportFilter) gatewaySender
.getGatewayTransportFilters().get(1)).name).isEqualTo("transportBean1");
assertThat(gatewaySenderConfigurer.beanNames.get(gatewaySender.getId()).toArray())
.isEqualTo(new String[] { "transportBean2", "transportBean1" });
Region region1 = (Region) this.applicationContext.getBean("Region1");
Region region2 = (Region) this.applicationContext.getBean("Region2");
assertThat(region1.getAttributes().getGatewaySenderIds())
.containsExactlyInAnyOrder("TestGatewaySender", "TestGatewaySender2");
assertThat(region2.getAttributes().getGatewaySenderIds())
.containsExactlyInAnyOrder("TestGatewaySender", "TestGatewaySender2");
}
}
@Test
public void annotationConfiguredMultipleGatewaySenders() {
this.applicationContext = newApplicationContext(BaseGatewaySenderTestConfiguration.class,
TestConfigurationWithMultipleGatewaySenderAnnotations.class);
Map<String, GatewaySender> beansOfType = this.applicationContext.getBeansOfType(GatewaySender.class);
assertThat(beansOfType.keySet().toArray()).containsExactlyInAnyOrder("TestGatewaySender", "TestGatewaySender2");
}
@Test
public void annotationConfiguredGatewaySender() {
this.applicationContext = newApplicationContext(BaseGatewaySenderTestConfiguration.class,
TestConfigurationWithAnnotations.class);
TestGatewaySenderConfigurer gatewaySenderConfigurer = this.applicationContext
.getBean(TestGatewaySenderConfigurer.class);
GatewaySender gatewaySender = this.applicationContext.getBean("TestGatewaySender", GatewaySender.class);
assertThat(gatewaySender.isManualStart()).isEqualTo(true);
assertThat(gatewaySender.getRemoteDSId()).isEqualTo(2);
assertThat(gatewaySender.getId()).isEqualTo("TestGatewaySender");
assertThat(gatewaySender.getDispatcherThreads()).isEqualTo(22);
assertThat(gatewaySender.isBatchConflationEnabled()).isEqualTo(true);
assertThat(gatewaySender.isParallel()).isEqualTo(true);
assertThat(gatewaySender.isPersistenceEnabled()).isEqualTo(false);
assertThat(gatewaySender.getDiskStoreName()).isEqualTo("someDiskStore");
assertThat(gatewaySender.getOrderPolicy()).isEqualTo(GatewaySender.OrderPolicy.PARTITION);
assertThat(gatewaySender.getGatewayEventFilters())
.containsExactlyInAnyOrder(this.applicationContext.getBean("SomeEventFilter", GatewayEventFilter.class));
assertThat(((TestGatewayEventSubstitutionFilter) gatewaySender.getGatewayEventSubstitutionFilter()).name)
.isEqualTo("SomeEventSubstitutionFilter");
assertThat(gatewaySender.getAlertThreshold()).isEqualTo(1234);
assertThat(gatewaySender.getBatchSize()).isEqualTo(100);
assertThat(gatewaySender.getBatchTimeInterval()).isEqualTo(2000);
assertThat(gatewaySender.getMaximumQueueMemory()).isEqualTo(400);
assertThat(gatewaySender.getSocketReadTimeout()).isEqualTo(4000);
assertThat(gatewaySender.getSocketBufferSize()).isEqualTo(16384);
assertThat(gatewaySender.getGatewayTransportFilters().size()).isEqualTo(2);
assertThat(gatewaySenderConfigurer.beanNames.get(gatewaySender.getId()).toArray())
.isEqualTo(new String[] { "transportBean2", "transportBean1" });
Region region1 = (Region) this.applicationContext.getBean("Region1");
Region region2 = (Region) this.applicationContext.getBean("Region2");
assertThat(region1.getAttributes().getGatewaySenderIds()).containsExactlyInAnyOrder("TestGatewaySender");
assertThat(region2.getAttributes().getGatewaySenderIds()).containsExactlyInAnyOrder("TestGatewaySender");
}
@Test
public void annotationConfigurationOfMultipleGatewaySendersWithOverrides() {
this.applicationContext = newApplicationContext(BaseGatewaySenderTestConfiguration.class,
TestConfigurationOfMultipleGatewaySenderAnnotationsWithOverrides.class);
TestGatewaySenderConfigurer gatewaySenderConfigurer = this.applicationContext
.getBean("gatewayConfigurer", TestGatewaySenderConfigurer.class);
Map<String, GatewaySender> beansOfType = this.applicationContext.getBeansOfType(GatewaySender.class);
assertThat(beansOfType.size()).isEqualTo(2);
String[] senders = new String[] { "TestGatewaySender", "TestGatewaySender2" };
assertThat(beansOfType.keySet().toArray()).containsExactly(senders);
for (String sender : senders) {
GatewaySender gatewaySender = this.applicationContext.getBean(sender, GatewaySender.class);
assertThat(gatewaySender.isManualStart()).isEqualTo(true);
assertThat(gatewaySender.getRemoteDSId()).isEqualTo(2);
assertThat(gatewaySender.getId()).isEqualTo(sender);
assertThat(gatewaySender.getDispatcherThreads()).isEqualTo(22);
assertThat(gatewaySender.isBatchConflationEnabled()).isEqualTo(true);
assertThat(gatewaySender.isParallel()).isEqualTo(true);
assertThat(gatewaySender.isPersistenceEnabled()).isEqualTo(true);
assertThat(gatewaySender.getDiskStoreName()).isEqualTo("someDiskStore");
assertThat(gatewaySender.getOrderPolicy()).isEqualTo(GatewaySender.OrderPolicy.PARTITION);
assertThat(gatewaySender.getGatewayEventSubstitutionFilter()).isNull();
assertThat(gatewaySender.getAlertThreshold()).isEqualTo(1234);
assertThat(gatewaySender.getBatchSize()).isEqualTo(1002);
assertThat(gatewaySender.getBatchTimeInterval()).isEqualTo(2000);
assertThat(gatewaySender.getMaximumQueueMemory()).isEqualTo(400);
assertThat(gatewaySender.getSocketReadTimeout()).isEqualTo(4000);
assertThat(gatewaySender.getSocketBufferSize()).isEqualTo(16384);
}
GatewaySender gatewaySender = this.applicationContext.getBean("TestGatewaySender", GatewaySender.class);
assertThat(gatewaySender.getGatewayTransportFilters().size()).isEqualTo(1);
assertThat(((GatewaySenderConfigurationTests.TestGatewayTransportFilter) gatewaySender
.getGatewayTransportFilters().get(0)).name).isEqualTo("transportBean1");
gatewaySender = this.applicationContext.getBean("TestGatewaySender2", GatewaySender.class);
assertThat(gatewaySender.getGatewayTransportFilters().size()).isEqualTo(2);
assertThat(gatewaySenderConfigurer.beanNames.get(gatewaySender.getId()).toArray())
.isEqualTo(new String[] { "transportBean2", "transportBean1" });
Region region1 = (Region) this.applicationContext.getBean("Region1");
Region region2 = (Region) this.applicationContext.getBean("Region2");
assertThat(region1.getAttributes().getGatewaySenderIds()).containsExactlyInAnyOrder("TestGatewaySender2");
assertThat(region2.getAttributes().getGatewaySenderIds())
.containsExactlyInAnyOrder("TestGatewaySender", "TestGatewaySender2");
}
private ConfigurableApplicationContext newApplicationContext(Class<?>... annotatedClasses) {
AnnotationConfigApplicationContext applicationContext = new AnnotationConfigApplicationContext(
annotatedClasses);
applicationContext.registerShutdownHook();
return applicationContext;
}
@EnableGatewaySenders(gatewaySenders = {
@EnableGatewaySender(name = "TestGatewaySender", manualStart = true, remoteDistributedSystemId = 2,
diskSynchronous = true, batchConflationEnabled = true, parallel = true, persistent = false,
diskStoreReference = "someDiskStore", orderPolicy = OrderPolicyType.PARTITION, eventFilters = "SomeEventFilter",
eventSubstitutionFilter = "SomeEventSubstitutionFilter", alertThreshold = 1234, batchSize = 100,
batchTimeInterval = 2000, dispatcherThreads = 22, maximumQueueMemory = 400, socketBufferSize = 16384,
socketReadTimeout = 4000, transportFilters = { "transportBean2", "transportBean1" },
regions = { "Region1", "Region2" })
})
static class TestConfigurationWithAnnotations {
@Bean("SomeEventSubstitutionFilter")
GatewayEventSubstitutionFilter createGatewayEventSubstitutionFilter() {
return new TestGatewayEventSubstitutionFilter("SomeEventSubstitutionFilter");
}
}
@EnableGatewaySenders(gatewaySenders = {
@EnableGatewaySender(name = "TestGatewaySender", manualStart = true, remoteDistributedSystemId = 2,
diskSynchronous = true, batchConflationEnabled = true, parallel = true, persistent = false,
diskStoreReference = "someDiskStore", orderPolicy = OrderPolicyType.PARTITION, alertThreshold = 1234, batchSize = 100,
eventFilters = "SomeEventFilter", batchTimeInterval = 2000, dispatcherThreads = 22, maximumQueueMemory = 400, socketBufferSize = 16384,
socketReadTimeout = 4000, regions = { "Region1", "Region2" }),
@EnableGatewaySender(name = "TestGatewaySender2", manualStart = true, remoteDistributedSystemId = 2,
diskSynchronous = true, batchConflationEnabled = true, parallel = true, persistent = false,
diskStoreReference = "someDiskStore", orderPolicy = OrderPolicyType.PARTITION, alertThreshold = 1234, batchSize = 100,
eventFilters = "SomeEventFilter", batchTimeInterval = 2000, dispatcherThreads = 22, maximumQueueMemory = 400, socketBufferSize = 16384,
socketReadTimeout = 4000, regions = { "Region1", "Region2" })
})
static class TestConfigurationWithMultipleGatewaySenderAnnotations {
@Bean("SomeEventSubstitutionFilter")
GatewayEventSubstitutionFilter createGatewayEventSubstitutionFilter() {
return new TestGatewayEventSubstitutionFilter("SomeEventSubstitutionFilter");
}
}
@EnableGatewaySenders(gatewaySenders = {
@EnableGatewaySender(name = "TestGatewaySender"),
@EnableGatewaySender(name = "TestGatewaySender2")
},
manualStart = true, remoteDistributedSystemId = 2,
diskSynchronous = false, batchConflationEnabled = true, parallel = true, persistent = true,
diskStoreReference = "someDiskStore", orderPolicy = OrderPolicyType.PARTITION, alertThreshold = 1234, batchSize = 1002,
eventFilters = "SomeEventFilter", batchTimeInterval = 2000, dispatcherThreads = 22, maximumQueueMemory = 400, socketBufferSize = 16384,
socketReadTimeout = 4000, regions = { "Region1", "Region2" },
transportFilters = { "transportBean2", "transportBean1" })
static class TestConfigurationOfMultipleGatewaySenderAnnotationsButWithDefaultsFromParent {
}
@EnableGatewaySenders(gatewaySenders = {
@EnableGatewaySender(name = "TestGatewaySender", transportFilters = "transportBean1", regions = "Region2"),
@EnableGatewaySender(name = "TestGatewaySender2")
},
manualStart = true, remoteDistributedSystemId = 2,
diskSynchronous = false, batchConflationEnabled = true, parallel = true, persistent = true,
diskStoreReference = "someDiskStore", orderPolicy = OrderPolicyType.PARTITION, alertThreshold = 1234, batchSize = 1002,
eventFilters = "SomeEventFilter", batchTimeInterval = 2000, dispatcherThreads = 22, maximumQueueMemory = 400, socketBufferSize = 16384,
socketReadTimeout = 4000, regions = { "Region1", "Region2" },
transportFilters = { "transportBean2", "transportBean1" })
static class TestConfigurationOfMultipleGatewaySenderAnnotationsWithOverrides {
}
private static class TestGatewaySenderConfigurer implements GatewaySenderConfigurer {
private final Map<String, List> beanNames = new TreeMap<>();
@Override
public void configure(String beanName, GatewaySenderFactoryBean bean) {
beanNames.put(beanName, bean.getTransportFilters().stream()
.map(transportFilter -> ((TestGatewayTransportFilter) transportFilter).name)
.collect(Collectors.toList()));
}
}
private static class TestGatewayEventSubstitutionFilter implements GatewayEventSubstitutionFilter {
private String name;
public TestGatewayEventSubstitutionFilter(String name) {
this.name = name;
}
@Override public Object getSubstituteValue(EntryEvent entryEvent) {
return null;
}
@Override public void close() {
}
}
private static class TestGatewayTransportFilter implements GatewayTransportFilter {
private String name;
public TestGatewayTransportFilter(String name) {
this.name = name;
}
@Override
public InputStream getInputStream(InputStream inputStream) {
return null;
}
@Override
public OutputStream getOutputStream(OutputStream outputStream) {
return null;
}
@Override public int hashCode() {
return name.hashCode();
}
@Override public boolean equals(Object obj) {
return this.name.equals(((TestGatewayTransportFilter) obj).name);
}
}
private static class TestGatewayEventFilter implements GatewayEventFilter {
private String name;
public TestGatewayEventFilter(String name) {
this.name = name;
}
@Override public boolean beforeEnqueue(GatewayQueueEvent gatewayQueueEvent) {
return false;
}
@Override public boolean beforeTransmit(GatewayQueueEvent gatewayQueueEvent) {
return false;
}
@Override public void afterAcknowledgement(GatewayQueueEvent gatewayQueueEvent) {
}
}
@PeerCacheApplication
@EnableGemFireMockObjects
static class BaseGatewaySenderTestConfiguration {
@Bean("Region1")
PartitionedRegionFactoryBean createRegion1(GemFireCache gemFireCache) {
return createRegion("Region1", gemFireCache);
}
@Bean("Region2")
PartitionedRegionFactoryBean createRegion2(GemFireCache gemFireCache) {
return createRegion("Region2", gemFireCache);
}
@Bean("gatewayConfigurer")
GatewaySenderConfigurer gatewaySenderConfigurer() {
return new TestGatewaySenderConfigurer();
}
@Bean("transportBean1")
GatewayTransportFilter createGatewayTransportBean1() {
return new TestGatewayTransportFilter("transportBean1");
}
@Bean("transportBean2")
GatewayTransportFilter createGatewayTransportBean2() {
return new TestGatewayTransportFilter("transportBean2");
}
@Bean("SomeEventFilter")
GatewayEventFilter createGatewayEventFilter() {
return new TestGatewayEventFilter("SomeEventFilter");
}
public PartitionedRegionFactoryBean createRegion(String name, GemFireCache gemFireCache) {
final PartitionedRegionFactoryBean regionFactoryBean = new PartitionedRegionFactoryBean();
regionFactoryBean.setCache(gemFireCache);
regionFactoryBean.setDataPolicy(DataPolicy.PARTITION);
regionFactoryBean.setName(name);
return regionFactoryBean;
}
}
}

View File

@@ -0,0 +1,250 @@
package org.springframework.data.gemfire.config.annotation;
import static org.assertj.core.api.Assertions.assertThat;
import java.io.InputStream;
import java.io.OutputStream;
import java.util.Map;
import java.util.Optional;
import org.apache.geode.cache.DataPolicy;
import org.apache.geode.cache.EntryEvent;
import org.apache.geode.cache.GemFireCache;
import org.apache.geode.cache.Region;
import org.apache.geode.cache.wan.GatewayEventFilter;
import org.apache.geode.cache.wan.GatewayEventSubstitutionFilter;
import org.apache.geode.cache.wan.GatewayQueueEvent;
import org.apache.geode.cache.wan.GatewaySender;
import org.apache.geode.cache.wan.GatewayTransportFilter;
import org.junit.After;
import org.junit.Before;
import org.junit.Test;
import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.context.annotation.AnnotationConfigApplicationContext;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.env.MutablePropertySources;
import org.springframework.core.env.PropertySource;
import org.springframework.data.gemfire.PartitionedRegionFactoryBean;
import org.springframework.data.gemfire.test.mock.annotation.EnableGemFireMockObjects;
import org.springframework.mock.env.MockPropertySource;
/**
* Tests for {@link EnableGatewaySenders} and {@link EnableGatewaySender}.
*
* @author Udo Kohlmeyer
* @see org.junit.Test
* @see org.mockito.Mockito
* @see GatewaySenderConfigurer
* @see GatewaySenderConfiguration
* @see org.apache.geode.cache.server.CacheServer
* @see org.springframework.context.annotation.Configuration
* @see org.springframework.test.context.ContextConfiguration
* @see org.springframework.test.context.junit4.SpringRunner
* @see org.springframework.data.gemfire.wan.GatewaySenderFactoryBean
* @since 2.2.0
*/
public class GatewaySenderConfigurerTests {
private ConfigurableApplicationContext applicationContext;
@Before
public void setup() {
}
@After
public void shutdown() {
Optional.ofNullable(this.applicationContext).ifPresent(ConfigurableApplicationContext::close);
}
private ConfigurableApplicationContext newApplicationContext(PropertySource<?> testPropertySource,
Class<?>... annotatedClasses) {
AnnotationConfigApplicationContext applicationContext = new AnnotationConfigApplicationContext();
MutablePropertySources propertySources = applicationContext.getEnvironment().getPropertySources();
propertySources.addFirst(testPropertySource);
applicationContext.registerShutdownHook();
applicationContext.register(annotatedClasses);
applicationContext.refresh();
return applicationContext;
}
@Test
public void annotationConfigurationOfMultipleGatewaySendersWithConfigurersAndProperties() {
MockPropertySource testPropertySource = new MockPropertySource();
testPropertySource.setProperty("spring.data.gemfire.gateway.sender.alert-threshold", 1234);
testPropertySource.setProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.batch-size", 1002);
testPropertySource.setProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.batch-size", 1002);
testPropertySource.setProperty("spring.data.gemfire.gateway.sender.batch-time-interval", 2000);
testPropertySource.setProperty("spring.data.gemfire.gateway.sender.maximum-queue-memory", 400);
testPropertySource
.setProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.socket-read-timeout", 4000);
testPropertySource
.setProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.socket-read-timeout", 4000);
testPropertySource.setProperty("spring.data.gemfire.gateway.sender.socket-buffer-size", 16384);
this.applicationContext = newApplicationContext(testPropertySource,
BaseGatewaySenderTestConfiguration.class,
TestTwoGatewaySenderConfigurersBasic.class);
Map<String, GatewaySender> beansOfType = this.applicationContext.getBeansOfType(GatewaySender.class);
assertThat(beansOfType.size()).isEqualTo(2);
String[] senders = new String[] { "TestGatewaySender", "TestGatewaySender2" };
assertThat(beansOfType.keySet().toArray()).containsExactly(senders);
for (String sender : senders) {
GatewaySender gatewaySender = this.applicationContext.getBean(sender, GatewaySender.class);
assertThat(gatewaySender.isManualStart()).isEqualTo(false);
assertThat(gatewaySender.getRemoteDSId()).isEqualTo(2);
assertThat(gatewaySender.getId()).isEqualTo(sender);
assertThat(gatewaySender.getDispatcherThreads()).isEqualTo(22);
assertThat(gatewaySender.isBatchConflationEnabled()).isEqualTo(false);
assertThat(gatewaySender.isParallel()).isEqualTo(true);
assertThat(gatewaySender.isPersistenceEnabled()).isEqualTo(false);
assertThat(gatewaySender.getOrderPolicy()).isEqualTo(GatewaySender.OrderPolicy.PARTITION);
assertThat(gatewaySender.getGatewayEventSubstitutionFilter()).isNull();
assertThat(gatewaySender.getAlertThreshold()).isEqualTo(1234);
assertThat(gatewaySender.getBatchSize()).isEqualTo(1002);
assertThat(gatewaySender.getBatchTimeInterval()).isEqualTo(2000);
assertThat(gatewaySender.getMaximumQueueMemory()).isEqualTo(400);
assertThat(gatewaySender.getSocketReadTimeout()).isEqualTo(4000);
assertThat(gatewaySender.getSocketBufferSize()).isEqualTo(16384);
assertThat(gatewaySender.getGatewayTransportFilters().size()).isEqualTo(0);
Region region1 = (Region) this.applicationContext.getBean("Region1");
Region region2 = (Region) this.applicationContext.getBean("Region2");
assertThat(region1.getAttributes().getGatewaySenderIds())
.containsExactlyInAnyOrder("TestGatewaySender", "TestGatewaySender2");
assertThat(region2.getAttributes().getGatewaySenderIds())
.containsExactlyInAnyOrder("TestGatewaySender2");
}
}
@PeerCacheApplication
@EnableGemFireMockObjects
static class BaseGatewaySenderTestConfiguration {
@Bean("Region1")
PartitionedRegionFactoryBean createRegion1(GemFireCache gemFireCache) {
return createRegion("Region1", gemFireCache);
}
@Bean("Region2")
PartitionedRegionFactoryBean createRegion2(GemFireCache gemFireCache) {
return createRegion("Region2", gemFireCache);
}
@Bean("transportBean1")
GatewayTransportFilter createGatewayTransportBean1() {
return new TestGatewayTransportFilter("transportBean1");
}
@Bean("transportBean2")
GatewayTransportFilter createGatewayTransportBean2() {
return new TestGatewayTransportFilter("transportBean2");
}
@Bean("SomeEventFilter")
GatewayEventFilter createGatewayEventFilter() {
return new TestGatewayEventFilter("SomeEventFilter");
}
public PartitionedRegionFactoryBean createRegion(String name, GemFireCache gemFireCache) {
final PartitionedRegionFactoryBean regionFactoryBean = new PartitionedRegionFactoryBean();
regionFactoryBean.setCache(gemFireCache);
regionFactoryBean.setDataPolicy(DataPolicy.PARTITION);
regionFactoryBean.setName(name);
return regionFactoryBean;
}
}
private static class TestGatewayEventSubstitutionFilter implements GatewayEventSubstitutionFilter {
private String name;
public TestGatewayEventSubstitutionFilter(String name) {
this.name = name;
}
@Override public Object getSubstituteValue(EntryEvent entryEvent) {
return null;
}
@Override public void close() {
}
}
private static class TestGatewayTransportFilter implements GatewayTransportFilter {
private String name;
public TestGatewayTransportFilter(String name) {
this.name = name;
}
@Override
public InputStream getInputStream(InputStream inputStream) {
return null;
}
@Override
public OutputStream getOutputStream(OutputStream outputStream) {
return null;
}
@Override public int hashCode() {
return name.hashCode();
}
@Override public boolean equals(Object obj) {
return this.name.equals(((TestGatewayTransportFilter) obj).name);
}
}
private static class TestGatewayEventFilter implements GatewayEventFilter {
private String name;
public TestGatewayEventFilter(String name) {
this.name = name;
}
@Override public boolean beforeEnqueue(GatewayQueueEvent gatewayQueueEvent) {
return false;
}
@Override public boolean beforeTransmit(GatewayQueueEvent gatewayQueueEvent) {
return false;
}
@Override public void afterAcknowledgement(GatewayQueueEvent gatewayQueueEvent) {
}
}
@Configuration
@EnableGatewaySenders(gatewaySenders = {
@EnableGatewaySender(name = "TestGatewaySender", regions = "Region1"),
@EnableGatewaySender(name = "TestGatewaySender2", regions = { "Region1", "Region2" })
})
static class TestTwoGatewaySenderConfigurersBasic {
@Bean
GatewaySenderConfigurer gatewaySenderConfigurer() {
return ((beanName, gatewaySenderFactoryBean) -> {
gatewaySenderFactoryBean.setRemoteDistributedSystemId(2);
gatewaySenderFactoryBean.setDispatcherThreads(22);
gatewaySenderFactoryBean.setParallel(true);
gatewaySenderFactoryBean.setOrderPolicy(GatewaySender.OrderPolicy.PARTITION);
});
}
}
}

View File

@@ -0,0 +1,671 @@
package org.springframework.data.gemfire.config.annotation;
import static org.assertj.core.api.Assertions.assertThat;
import java.io.InputStream;
import java.io.OutputStream;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.TreeMap;
import java.util.stream.Collectors;
import org.apache.geode.cache.DataPolicy;
import org.apache.geode.cache.EntryEvent;
import org.apache.geode.cache.GemFireCache;
import org.apache.geode.cache.Region;
import org.apache.geode.cache.wan.GatewayEventFilter;
import org.apache.geode.cache.wan.GatewayEventSubstitutionFilter;
import org.apache.geode.cache.wan.GatewayQueueEvent;
import org.apache.geode.cache.wan.GatewaySender;
import org.apache.geode.cache.wan.GatewayTransportFilter;
import org.apache.geode.distributed.internal.DistributionAdvisor;
import org.apache.geode.internal.cache.EntryEventImpl;
import org.apache.geode.internal.cache.wan.AbstractGatewaySender;
import org.junit.After;
import org.junit.Before;
import org.junit.Test;
import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.context.annotation.AnnotationConfigApplicationContext;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.env.MutablePropertySources;
import org.springframework.core.env.PropertySource;
import org.springframework.data.gemfire.PartitionedRegionFactoryBean;
import org.springframework.data.gemfire.test.mock.annotation.EnableGemFireMockObjects;
import org.springframework.data.gemfire.wan.GatewaySenderFactoryBean;
import org.springframework.data.gemfire.wan.OrderPolicyType;
import org.springframework.mock.env.MockPropertySource;
/**
* Tests for {@link EnableGatewaySenders} and {@link EnableGatewaySender} to test the configuration of {@link GatewaySender} using properties.
*
* @author Udo Kohlmeyer
* @see Test
* @see Configuration
* @see org.mockito.Mockito
* @see GatewaySenderConfigurer
* @see GatewaySenderConfiguration
* @see org.apache.geode.cache.server.CacheServer
* @see org.springframework.test.context.ContextConfiguration
* @see org.springframework.test.context.junit4.SpringRunner
* @see org.springframework.data.gemfire.wan.GatewayReceiverFactoryBean
* @since 2.2.0
*/
public class GatewaySenderPropertiesTests {
private ConfigurableApplicationContext applicationContext;
@Before
public void setup() {
}
@After
public void shutdown() {
Optional.ofNullable(this.applicationContext).ifPresent(ConfigurableApplicationContext::close);
}
@Test
public void gatewayReceiverPropertiesConfigurationOnMultipleChildren() {
MockPropertySource testPropertySource = new MockPropertySource()
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.manual-start", true)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.remote-distributed-system-id", 2)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.disk-synchronous", true)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.batch-conflation-enabled", true)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.parallel", true)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.persistent", false)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.order-policy", "PARTITION")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.event-substitution-filter",
"SomeEventSubstitutionFilter")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.alert-threshold", 1234)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.batch-size", 100)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.batch-time-interval", 2000)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.dispatcher-threads", 22)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.maximum-queue-memory", 400)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.socket-buffer-size", 16384)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.socket-read-timeout", 4000)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.disk-store-reference", "someDiskStore")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.event-filters", "SomeEventFilter")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.transport-filters",
"transportBean2, transportBean1")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.regions", "Region1,Region2")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.manual-start", false)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.remote-distributed-system-id", 3)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.disk-synchronous", false)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.batch-conflation-enabled", false)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.parallel", false)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.persistent", true)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.order-policy", "KEY")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.event-substitution-filter",
"SomeEventSubstitutionFilter")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.alert-threshold", 4321)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.batch-size", 1000)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.batch-time-interval", 20000)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.dispatcher-threads", 2200)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.maximum-queue-memory", 40000)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.socket-buffer-size", 1638400)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.socket-read-timeout", 400000)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.disk-store-reference", "someDiskStore")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.event-filters", "SomeEventFilter")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.transport-filters",
"transportBean1")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.regions", "Region1");
this.applicationContext = newApplicationContext(testPropertySource, BaseGatewaySenderTestConfiguration.class,
TestConfigurationWithPropertiesMultipleGatewaySenders.class);
TestGatewaySenderConfigurer gatewaySenderConfigurer = this.applicationContext
.getBean(TestGatewaySenderConfigurer.class);
GatewaySender gatewaySender = this.applicationContext.getBean("TestGatewaySender", GatewaySender.class);
assertThat(gatewaySender.isManualStart()).isEqualTo(true);
assertThat(gatewaySender.getRemoteDSId()).isEqualTo(2);
assertThat(gatewaySender.getId()).isEqualTo("TestGatewaySender");
assertThat(gatewaySender.getDispatcherThreads()).isEqualTo(22);
assertThat(gatewaySender.isBatchConflationEnabled()).isEqualTo(true);
assertThat(gatewaySender.isParallel()).isEqualTo(true);
assertThat(gatewaySender.isPersistenceEnabled()).isEqualTo(false);
assertThat(gatewaySender.getDiskStoreName()).isEqualTo("someDiskStore");
assertThat(gatewaySender.getOrderPolicy()).isEqualTo(GatewaySender.OrderPolicy.PARTITION);
assertThat(((TestGatewayEventSubstitutionFilter) gatewaySender.getGatewayEventSubstitutionFilter()).name)
.isEqualTo("SomeEventSubstitutionFilter");
assertThat(gatewaySender.getAlertThreshold()).isEqualTo(1234);
assertThat(gatewaySender.getBatchSize()).isEqualTo(100);
assertThat(gatewaySender.getBatchTimeInterval()).isEqualTo(2000);
assertThat(gatewaySender.getMaximumQueueMemory()).isEqualTo(400);
assertThat(gatewaySender.getSocketReadTimeout()).isEqualTo(4000);
assertThat(gatewaySender.getSocketBufferSize()).isEqualTo(16384);
assertThat(gatewaySender.getGatewayTransportFilters().size()).isEqualTo(2);
assertThat(gatewaySenderConfigurer.beanNames.get(gatewaySender.getId()).toArray())
.isEqualTo(new String[] { "transportBean2", "transportBean1" });
gatewaySender = this.applicationContext.getBean("TestGatewaySender2", GatewaySender.class);
assertThat(gatewaySender.isManualStart()).isEqualTo(false);
assertThat(gatewaySender.getRemoteDSId()).isEqualTo(3);
assertThat(gatewaySender.getId()).isEqualTo("TestGatewaySender2");
assertThat(gatewaySender.getDispatcherThreads()).isEqualTo(2200);
assertThat(gatewaySender.isBatchConflationEnabled()).isEqualTo(false);
assertThat(gatewaySender.isParallel()).isEqualTo(false);
assertThat(gatewaySender.isPersistenceEnabled()).isEqualTo(true);
assertThat(gatewaySender.getDiskStoreName()).isEqualTo("someDiskStore");
assertThat(gatewaySender.getOrderPolicy()).isEqualTo(GatewaySender.OrderPolicy.KEY);
assertThat(((TestGatewayEventSubstitutionFilter) gatewaySender.getGatewayEventSubstitutionFilter()).name)
.isEqualTo("SomeEventSubstitutionFilter");
assertThat(gatewaySender.getAlertThreshold()).isEqualTo(4321);
assertThat(gatewaySender.getBatchSize()).isEqualTo(1000);
assertThat(gatewaySender.getBatchTimeInterval()).isEqualTo(20000);
assertThat(gatewaySender.getMaximumQueueMemory()).isEqualTo(40000);
assertThat(gatewaySender.getSocketReadTimeout()).isEqualTo(400000);
assertThat(gatewaySender.getSocketBufferSize()).isEqualTo(1638400);
assertThat(gatewaySender.getGatewayTransportFilters().size()).isEqualTo(1);
assertThat(gatewaySenderConfigurer.beanNames.get(gatewaySender.getId()).toArray())
.isEqualTo(new String[] { "transportBean1" });
Region region1 = (Region) this.applicationContext.getBean("Region1");
Region region2 = (Region) this.applicationContext.getBean("Region2");
assertThat(region1.getAttributes().getGatewaySenderIds())
.containsExactlyInAnyOrder("TestGatewaySender", "TestGatewaySender2");
assertThat(region2.getAttributes().getGatewaySenderIds())
.containsExactlyInAnyOrder("TestGatewaySender");
}
@Test
public void gatewayReceiverPropertiesConfigurationOnMultipleChildrenAndAnnotations() {
MockPropertySource testPropertySource = new MockPropertySource()
.withProperty("spring.data.gemfire.gateway.sender.socket-read-timeout", 4000)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.manual-start", true)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.remote-distributed-system-id", 2)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.disk-synchronous", true)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.batch-conflation-enabled", true)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.parallel", true)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.persistent", false)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.order-policy", "THREAD")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.event-substitution-filter",
"SomeEventSubstitutionFilter")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.alert-threshold", 1234)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.batch-size", 1020)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.batch-time-interval", 2300)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.dispatcher-threads", 22)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.maximum-queue-memory", 400)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.socket-buffer-size", 16384)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.disk-store-reference", "someDiskStore")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.event-filters", "SomeEventFilter")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.transport-filters",
"transportBean2, transportBean1")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.regions", "Region1,Region2")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.manual-start", false)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.remote-distributed-system-id", 3)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.disk-synchronous", false)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.batch-conflation-enabled", false)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.parallel", false)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.persistent", true)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.order-policy", "KEY")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.event-substitution-filter",
"SomeEventSubstitutionFilter")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.alert-threshold", 4321)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.batch-size", 1000)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.batch-time-interval", 20000)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.dispatcher-threads", 2200)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.maximum-queue-memory", 40000)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.socket-buffer-size", 1638400)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.disk-store-reference", "someDiskStore")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.event-filters", "SomeEventFilter")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.transport-filters",
"transportBean1")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender2.regions", "");
this.applicationContext = newApplicationContext(testPropertySource, BaseGatewaySenderTestConfiguration.class,
TestConfigurationWithMultipleGatewaySenderAnnotations.class);
TestGatewaySenderConfigurer gatewaySenderConfigurer = this.applicationContext
.getBean(TestGatewaySenderConfigurer.class);
GatewaySender gatewaySender = this.applicationContext.getBean("TestGatewaySender", GatewaySender.class);
assertThat(gatewaySender.isManualStart()).isEqualTo(true);
assertThat(gatewaySender.getRemoteDSId()).isEqualTo(2);
assertThat(gatewaySender.getId()).isEqualTo("TestGatewaySender");
assertThat(gatewaySender.getDispatcherThreads()).isEqualTo(22);
assertThat(gatewaySender.isBatchConflationEnabled()).isEqualTo(true);
assertThat(gatewaySender.isParallel()).isEqualTo(true);
assertThat(gatewaySender.isPersistenceEnabled()).isEqualTo(false);
assertThat(gatewaySender.getDiskStoreName()).isEqualTo("someDiskStore");
assertThat(gatewaySender.getOrderPolicy()).isEqualTo(GatewaySender.OrderPolicy.THREAD);
assertThat(((TestGatewayEventSubstitutionFilter) gatewaySender.getGatewayEventSubstitutionFilter()).name)
.isEqualTo("SomeEventSubstitutionFilter");
assertThat(gatewaySender.getAlertThreshold()).isEqualTo(1234);
assertThat(gatewaySender.getBatchSize()).isEqualTo(1020);
assertThat(gatewaySender.getBatchTimeInterval()).isEqualTo(2300);
assertThat(gatewaySender.getMaximumQueueMemory()).isEqualTo(400);
assertThat(gatewaySender.getSocketReadTimeout()).isEqualTo(4000);
assertThat(gatewaySender.getSocketBufferSize()).isEqualTo(16384);
assertThat(gatewaySender.getGatewayTransportFilters().size()).isEqualTo(2);
assertThat(gatewaySenderConfigurer.beanNames.get(gatewaySender.getId()).toArray())
.isEqualTo(new String[] { "transportBean2", "transportBean1" });
gatewaySender = this.applicationContext.getBean("TestGatewaySender2", GatewaySender.class);
assertThat(gatewaySender.isManualStart()).isEqualTo(false);
assertThat(gatewaySender.getRemoteDSId()).isEqualTo(3);
assertThat(gatewaySender.getId()).isEqualTo("TestGatewaySender2");
assertThat(gatewaySender.getDispatcherThreads()).isEqualTo(2200);
assertThat(gatewaySender.isBatchConflationEnabled()).isEqualTo(false);
assertThat(gatewaySender.isParallel()).isEqualTo(false);
assertThat(gatewaySender.isPersistenceEnabled()).isEqualTo(true);
assertThat(gatewaySender.getDiskStoreName()).isEqualTo("someDiskStore");
assertThat(gatewaySender.getOrderPolicy()).isEqualTo(GatewaySender.OrderPolicy.KEY);
assertThat(((TestGatewayEventSubstitutionFilter) gatewaySender.getGatewayEventSubstitutionFilter()).name)
.isEqualTo("SomeEventSubstitutionFilter");
assertThat(gatewaySender.getAlertThreshold()).isEqualTo(4321);
assertThat(gatewaySender.getBatchSize()).isEqualTo(1000);
assertThat(gatewaySender.getBatchTimeInterval()).isEqualTo(20000);
assertThat(gatewaySender.getMaximumQueueMemory()).isEqualTo(40000);
assertThat(gatewaySender.getSocketReadTimeout()).isEqualTo(4000);
assertThat(gatewaySender.getSocketBufferSize()).isEqualTo(1638400);
assertThat(gatewaySender.getGatewayTransportFilters().size()).isEqualTo(1);
assertThat(gatewaySenderConfigurer.beanNames.get(gatewaySender.getId()).toArray())
.isEqualTo(new String[] { "transportBean1" });
Region region1 = (Region) this.applicationContext.getBean("Region1");
Region region2 = (Region) this.applicationContext.getBean("Region2");
assertThat(region1.getAttributes().getGatewaySenderIds())
.containsExactlyInAnyOrder("TestGatewaySender", "TestGatewaySender2");
assertThat(region2.getAttributes().getGatewaySenderIds())
.containsExactlyInAnyOrder("TestGatewaySender", "TestGatewaySender2");
}
@Test
public void gatewayReceiverPropertiesConfigurationOnChild() {
MockPropertySource testPropertySource = new MockPropertySource()
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.manual-start", true)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.remote-distributed-system-id", 2)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.disk-synchronous", true)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.batch-conflation-enabled", true)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.parallel", true)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.persistent", false)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.order-policy", "PARTITION")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.event-substitution-filter",
"SomeEventSubstitutionFilter")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.alert-threshold", 1234)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.batch-size", 100)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.batch-time-interval", 2000)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.dispatcher-threads", 22)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.maximum-queue-memory", 400)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.socket-buffer-size", 16384)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.socket-read-timeout", 4000)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.disk-store-reference", "someDiskStore")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.event-filters", "SomeEventFilter")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.transport-filters",
"transportBean2, transportBean1")
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.regions", "Region1,Region2");
this.applicationContext = newApplicationContext(testPropertySource, BaseGatewaySenderTestConfiguration.class,
TestConfigurationWithProperties.class);
TestGatewaySenderConfigurer gatewaySenderConfigurer = this.applicationContext
.getBean(TestGatewaySenderConfigurer.class);
GatewaySender gatewaySender = this.applicationContext.getBean("TestGatewaySender", GatewaySender.class);
assertThat(gatewaySender.isManualStart()).isEqualTo(true);
assertThat(gatewaySender.getRemoteDSId()).isEqualTo(2);
assertThat(gatewaySender.getId()).isEqualTo("TestGatewaySender");
assertThat(gatewaySender.getDispatcherThreads()).isEqualTo(22);
assertThat(gatewaySender.isBatchConflationEnabled()).isEqualTo(true);
assertThat(gatewaySender.isParallel()).isEqualTo(true);
assertThat(gatewaySender.isPersistenceEnabled()).isEqualTo(false);
assertThat(gatewaySender.getDiskStoreName()).isEqualTo("someDiskStore");
assertThat(gatewaySender.getOrderPolicy()).isEqualTo(GatewaySender.OrderPolicy.PARTITION);
assertThat(((TestGatewayEventSubstitutionFilter) gatewaySender.getGatewayEventSubstitutionFilter()).name)
.isEqualTo("SomeEventSubstitutionFilter");
assertThat(gatewaySender.getAlertThreshold()).isEqualTo(1234);
assertThat(gatewaySender.getBatchSize()).isEqualTo(100);
assertThat(gatewaySender.getBatchTimeInterval()).isEqualTo(2000);
assertThat(gatewaySender.getMaximumQueueMemory()).isEqualTo(400);
assertThat(gatewaySender.getSocketReadTimeout()).isEqualTo(4000);
assertThat(gatewaySender.getSocketBufferSize()).isEqualTo(16384);
assertThat(gatewaySender.getGatewayTransportFilters().size()).isEqualTo(2);
assertThat(gatewaySenderConfigurer.beanNames.get(gatewaySender.getId()).toArray())
.isEqualTo(new String[] { "transportBean2", "transportBean1" });
Region region1 = (Region) this.applicationContext.getBean("Region1");
Region region2 = (Region) this.applicationContext.getBean("Region2");
assertThat(region1.getAttributes().getGatewaySenderIds())
.containsExactlyInAnyOrder("TestGatewaySender");
assertThat(region2.getAttributes().getGatewaySenderIds())
.containsExactlyInAnyOrder("TestGatewaySender");
}
@Test
public void gatewayReceiverPropertiesConfigurationOnParent() {
MockPropertySource testPropertySource = new MockPropertySource()
.withProperty("spring.data.gemfire.gateway.sender.manual-start", true)
.withProperty("spring.data.gemfire.gateway.sender.remote-distributed-system-id", 2)
.withProperty("spring.data.gemfire.gateway.sender.disk-synchronous", true)
.withProperty("spring.data.gemfire.gateway.sender.batch-conflation-enabled", true)
.withProperty("spring.data.gemfire.gateway.sender.parallel", true)
.withProperty("spring.data.gemfire.gateway.sender.persistent", false)
.withProperty("spring.data.gemfire.gateway.sender.order-policy", "PARTITION")
.withProperty("spring.data.gemfire.gateway.sender.event-substitution-filter",
"SomeEventSubstitutionFilter")
.withProperty("spring.data.gemfire.gateway.sender.alert-threshold", 1234)
.withProperty("spring.data.gemfire.gateway.sender.batch-size", 100)
.withProperty("spring.data.gemfire.gateway.sender.batch-time-interval", 2000)
.withProperty("spring.data.gemfire.gateway.sender.dispatcher-threads", 22)
.withProperty("spring.data.gemfire.gateway.sender.maximum-queue-memory", 400)
.withProperty("spring.data.gemfire.gateway.sender.socket-buffer-size", 16384)
.withProperty("spring.data.gemfire.gateway.sender.socket-read-timeout", 4000)
.withProperty("spring.data.gemfire.gateway.sender.disk-store-reference", "someDiskStore")
.withProperty("spring.data.gemfire.gateway.sender.event-filters", "SomeEventFilter")
.withProperty("spring.data.gemfire.gateway.sender.transport-filters",
"transportBean2, transportBean1")
.withProperty("spring.data.gemfire.gateway.sender.regions", "Region1,Region2");
this.applicationContext = newApplicationContext(testPropertySource, BaseGatewaySenderTestConfiguration.class,
TestConfigurationWithProperties.class);
TestGatewaySenderConfigurer gatewaySenderConfigurer = this.applicationContext
.getBean(TestGatewaySenderConfigurer.class);
GatewaySender gatewaySender = this.applicationContext.getBean("TestGatewaySender", GatewaySender.class);
assertThat(gatewaySender.isManualStart()).isEqualTo(true);
assertThat(gatewaySender.getRemoteDSId()).isEqualTo(2);
assertThat(gatewaySender.getId()).isEqualTo("TestGatewaySender");
assertThat(gatewaySender.getDispatcherThreads()).isEqualTo(22);
assertThat(gatewaySender.isBatchConflationEnabled()).isEqualTo(true);
assertThat(gatewaySender.isParallel()).isEqualTo(true);
assertThat(gatewaySender.isPersistenceEnabled()).isEqualTo(false);
assertThat(gatewaySender.getDiskStoreName()).isEqualTo("someDiskStore");
assertThat(gatewaySender.getOrderPolicy()).isEqualTo(GatewaySender.OrderPolicy.PARTITION);
assertThat(((TestGatewayEventSubstitutionFilter) gatewaySender.getGatewayEventSubstitutionFilter()).name)
.isEqualTo("SomeEventSubstitutionFilter");
assertThat(gatewaySender.getAlertThreshold()).isEqualTo(1234);
assertThat(gatewaySender.getBatchSize()).isEqualTo(100);
assertThat(gatewaySender.getBatchTimeInterval()).isEqualTo(2000);
assertThat(gatewaySender.getMaximumQueueMemory()).isEqualTo(400);
assertThat(gatewaySender.getSocketReadTimeout()).isEqualTo(4000);
assertThat(gatewaySender.getSocketBufferSize()).isEqualTo(16384);
assertThat(gatewaySender.getGatewayTransportFilters().size()).isEqualTo(2);
assertThat(gatewaySenderConfigurer.beanNames.get(gatewaySender.getId()).toArray())
.isEqualTo(new String[] { "transportBean2", "transportBean1" });
Region region1 = (Region) this.applicationContext.getBean("Region1");
Region region2 = (Region) this.applicationContext.getBean("Region2");
assertThat(region1.getAttributes().getGatewaySenderIds())
.containsExactlyInAnyOrder("TestGatewaySender");
assertThat(region2.getAttributes().getGatewaySenderIds())
.containsExactlyInAnyOrder("TestGatewaySender");
}
@Test
public void gatewayReceiverPropertiesConfigurationOnParentWithChildOverride() {
MockPropertySource testPropertySource = new MockPropertySource()
.withProperty("spring.data.gemfire.gateway.sender.manual-start", true)
.withProperty("spring.data.gemfire.gateway.sender.TestGatewaySender.manual-start", false)
.withProperty("spring.data.gemfire.gateway.sender.remote-distributed-system-id", 2)
.withProperty("spring.data.gemfire.gateway.sender.disk-synchronous", true)
.withProperty("spring.data.gemfire.gateway.sender.batch-conflation-enabled", true)
.withProperty("spring.data.gemfire.gateway.sender.parallel", true)
.withProperty("spring.data.gemfire.gateway.sender.persistent", false)
.withProperty("spring.data.gemfire.gateway.sender.order-policy", "PARTITION")
.withProperty("spring.data.gemfire.gateway.sender.event-substitution-filter",
"SomeEventSubstitutionFilter")
.withProperty("spring.data.gemfire.gateway.sender.alert-threshold", 1234)
.withProperty("spring.data.gemfire.gateway.sender.batch-size", 100)
.withProperty("spring.data.gemfire.gateway.sender.batch-time-interval", 2000)
.withProperty("spring.data.gemfire.gateway.sender.dispatcher-threads", 22)
.withProperty("spring.data.gemfire.gateway.sender.maximum-queue-memory", 400)
.withProperty("spring.data.gemfire.gateway.sender.socket-buffer-size", 16384)
.withProperty("spring.data.gemfire.gateway.sender.socket-read-timeout", 4000)
.withProperty("spring.data.gemfire.gateway.sender.disk-store-reference", "someDiskStore")
.withProperty("spring.data.gemfire.gateway.sender.event-filters", "SomeEventFilter")
.withProperty("spring.data.gemfire.gateway.sender.transport-filters",
"transportBean2, transportBean1")
.withProperty("spring.data.gemfire.gateway.sender.regions", "Region1,Region2");
this.applicationContext = newApplicationContext(testPropertySource, BaseGatewaySenderTestConfiguration.class,
TestConfigurationWithProperties.class);
TestGatewaySenderConfigurer gatewaySenderConfigurer = this.applicationContext
.getBean(TestGatewaySenderConfigurer.class);
GatewaySender gatewaySender = this.applicationContext.getBean("TestGatewaySender", GatewaySender.class);
assertThat(gatewaySender.isManualStart()).isEqualTo(false);
assertThat(gatewaySender.getRemoteDSId()).isEqualTo(2);
assertThat(gatewaySender.getId()).isEqualTo("TestGatewaySender");
assertThat(gatewaySender.getDispatcherThreads()).isEqualTo(22);
assertThat(gatewaySender.isBatchConflationEnabled()).isEqualTo(true);
assertThat(gatewaySender.isParallel()).isEqualTo(true);
assertThat(gatewaySender.isPersistenceEnabled()).isEqualTo(false);
assertThat(gatewaySender.getDiskStoreName()).isEqualTo("someDiskStore");
assertThat(gatewaySender.getOrderPolicy()).isEqualTo(GatewaySender.OrderPolicy.PARTITION);
assertThat(((TestGatewayEventSubstitutionFilter) gatewaySender.getGatewayEventSubstitutionFilter()).name)
.isEqualTo("SomeEventSubstitutionFilter");
assertThat(gatewaySender.getAlertThreshold()).isEqualTo(1234);
assertThat(gatewaySender.getBatchSize()).isEqualTo(100);
assertThat(gatewaySender.getBatchTimeInterval()).isEqualTo(2000);
assertThat(gatewaySender.getMaximumQueueMemory()).isEqualTo(400);
assertThat(gatewaySender.getSocketReadTimeout()).isEqualTo(4000);
assertThat(gatewaySender.getSocketBufferSize()).isEqualTo(16384);
assertThat(gatewaySender.getGatewayTransportFilters().size()).isEqualTo(2);
assertThat(gatewaySenderConfigurer.beanNames.get(gatewaySender.getId()).toArray())
.isEqualTo(new String[] { "transportBean2", "transportBean1" });
Region region1 = (Region) this.applicationContext.getBean("Region1");
Region region2 = (Region) this.applicationContext.getBean("Region2");
assertThat(region1.getAttributes().getGatewaySenderIds())
.containsExactlyInAnyOrder("TestGatewaySender");
assertThat(region2.getAttributes().getGatewaySenderIds())
.containsExactlyInAnyOrder("TestGatewaySender");
}
private ConfigurableApplicationContext newApplicationContext(PropertySource<?> testPropertySource,
Class<?>... annotatedClasses) {
AnnotationConfigApplicationContext applicationContext = new AnnotationConfigApplicationContext();
MutablePropertySources propertySources = applicationContext.getEnvironment().getPropertySources();
propertySources.addFirst(testPropertySource);
applicationContext.registerShutdownHook();
applicationContext.register(annotatedClasses);
applicationContext.refresh();
return applicationContext;
}
@EnableGatewaySenders(gatewaySenders = {
@EnableGatewaySender(name = "TestGatewaySender", manualStart = true, remoteDistributedSystemId = 2,
diskSynchronous = true, batchConflationEnabled = true, parallel = true, persistent = false,
diskStoreReference = "someDiskStore", orderPolicy = OrderPolicyType.PARTITION, alertThreshold = 1234, batchSize = 100,
batchTimeInterval = 2000, dispatcherThreads = 22, maximumQueueMemory = 400, socketBufferSize = 16384,
socketReadTimeout = 4000, regions = { "Region1", "Region2" }),
@EnableGatewaySender(name = "TestGatewaySender2", manualStart = true, remoteDistributedSystemId = 2,
diskSynchronous = true, batchConflationEnabled = true, parallel = true, persistent = false,
diskStoreReference = "someDiskStore", orderPolicy = OrderPolicyType.PARTITION, alertThreshold = 1234, batchSize = 100,
batchTimeInterval = 2000, dispatcherThreads = 22, maximumQueueMemory = 400, socketBufferSize = 16384,
socketReadTimeout = 4000, regions = { "Region1", "Region2" })
})
static class TestConfigurationWithMultipleGatewaySenderAnnotations {
}
@EnableGatewaySender(name = "TestGatewaySender")
static class TestConfigurationWithProperties {
@Bean("gatewayConfigurer")
GatewaySenderConfigurer gatewaySenderConfigurer() {
return new TestGatewaySenderConfigurer();
}
}
@EnableGatewaySenders(gatewaySenders = {
@EnableGatewaySender(name = "TestGatewaySender"),
@EnableGatewaySender(name = "TestGatewaySender2")
})
static class TestConfigurationWithPropertiesMultipleGatewaySenders {
}
private static class TestGatewaySenderConfigurer implements GatewaySenderConfigurer {
private final Map<String, List> beanNames = new TreeMap<>();
@Override
public void configure(String beanName, GatewaySenderFactoryBean bean) {
beanNames.put(beanName, bean.getTransportFilters().stream()
.map(transportFilter -> ((TestGatewayTransportFilter) transportFilter).name)
.collect(Collectors.toList()));
}
}
private static class TestGatewayEventSubstitutionFilter implements GatewayEventSubstitutionFilter {
private String name;
public TestGatewayEventSubstitutionFilter(String name) {
this.name = name;
}
@Override public Object getSubstituteValue(EntryEvent entryEvent) {
return null;
}
@Override public void close() {
}
}
private static class TestGatewayEventFilter implements GatewayEventFilter {
private String name;
public TestGatewayEventFilter(String name) {
this.name = name;
}
@Override public boolean beforeEnqueue(GatewayQueueEvent gatewayQueueEvent) {
return false;
}
@Override public boolean beforeTransmit(GatewayQueueEvent gatewayQueueEvent) {
return false;
}
@Override public void afterAcknowledgement(GatewayQueueEvent gatewayQueueEvent) {
}
}
private static class TestGatewayTransportFilter implements GatewayTransportFilter {
private String name;
public TestGatewayTransportFilter(String name) {
this.name = name;
}
@Override
public InputStream getInputStream(InputStream inputStream) {
return null;
}
@Override
public OutputStream getOutputStream(OutputStream outputStream) {
return null;
}
@Override public int hashCode() {
return name.hashCode();
}
@Override public boolean equals(Object obj) {
return this.name.equals(((TestGatewayTransportFilter) obj).name);
}
}
private static class TestGatewaySender extends AbstractGatewaySender implements GatewaySender {
@Override public void start() {
}
@Override public void stop() {
}
@Override public void setModifiedEventId(EntryEventImpl entryEvent) {
}
@Override public void fillInProfile(DistributionAdvisor.Profile profile) {
}
}
@PeerCacheApplication
@EnableGemFireMockObjects
static class BaseGatewaySenderTestConfiguration {
@Bean("Region1")
PartitionedRegionFactoryBean createRegion1(GemFireCache gemFireCache) {
return createRegion("Region1", gemFireCache);
}
@Bean("Region2")
PartitionedRegionFactoryBean createRegion2(GemFireCache gemFireCache) {
return createRegion("Region2", gemFireCache);
}
@Bean("gatewayConfigurer")
GatewaySenderConfigurer gatewaySenderConfigurer() {
return new TestGatewaySenderConfigurer();
}
@Bean("transportBean1")
GatewayTransportFilter createGatewayTransportBean1() {
return new TestGatewayTransportFilter("transportBean1");
}
@Bean("transportBean2")
GatewayTransportFilter createGatewayTransportBean2() {
return new TestGatewayTransportFilter("transportBean2");
}
@Bean("SomeEventFilter")
GatewayEventFilter createGatewayEventFilter() {
return new TestGatewayEventFilter("SomeEventFilter");
}
@Bean("SomeEventSubstitutionFilter")
GatewayEventSubstitutionFilter createGatewayEventSubstitutionFilter() {
return new TestGatewayEventSubstitutionFilter("SomeEventSubstitutionFilter");
}
public PartitionedRegionFactoryBean createRegion(String name, GemFireCache gemFireCache) {
final PartitionedRegionFactoryBean regionFactoryBean = new PartitionedRegionFactoryBean();
regionFactoryBean.setCache(gemFireCache);
regionFactoryBean.setDataPolicy(DataPolicy.PARTITION);
regionFactoryBean.setName(name);
return regionFactoryBean;
}
}
}

View File

@@ -16,6 +16,8 @@
package org.springframework.data.gemfire.repository.sample;
import java.util.Objects;
import org.springframework.data.annotation.Id;
import org.springframework.data.gemfire.mapping.annotation.Region;
import org.springframework.data.gemfire.mapping.annotation.ReplicateRegion;
@@ -58,8 +60,24 @@ public class Customer {
return String.format("%1$s %2$s", getFirstName(), getLastName());
}
public void setId(Long id) {
this.id = id;
}
public Long getId() {
return id;
}
public String getFirstName() {
return firstName;
}
public String getLastName() {
return lastName;
}
protected static boolean equalsIgnoreNull(final Object obj1, final Object obj2) {
return (obj1 == null ? obj2 == null : obj1.equals(obj2));
return (Objects.equals(obj1, obj2));
}
@Override

View File

@@ -111,11 +111,11 @@ public class StubGatewaySenderFactory implements GatewaySenderFactory {
}
});
doAnswer(new Answer<Void>() {
public Void answer(InvocationOnMock invocation) {
public Void answer(InvocationOnMock invocation) {
running = true;
return null;
}
}).when(gatewaySender).start();
}
}).when(gatewaySender).start();
return gatewaySender;
}
@@ -222,7 +222,8 @@ public class StubGatewaySenderFactory implements GatewaySenderFactory {
}
@Override
public GatewaySenderFactory setGatewayEventSubstitutionFilter(final GatewayEventSubstitutionFilter gatewayEventSubstitutionFilter) {
public GatewaySenderFactory setGatewayEventSubstitutionFilter(
final GatewayEventSubstitutionFilter gatewayEventSubstitutionFilter) {
this.gatewayEventSubstitutionFilter = gatewayEventSubstitutionFilter;
return this;
}

View File

@@ -218,7 +218,7 @@ public class GatewaySenderFactoryBeanTest {
factoryBean.setName("g6");
factoryBean.setRemoteDistributedSystemId(51);
factoryBean.setPersistent(false);
factoryBean.setDiskStoreRef("queueOverflowDiskStore");
factoryBean.setDiskStoreReference("queueOverflowDiskStore");
factoryBean.doInit();
verifyExpectations(factoryBean, mockGatewaySenderFactory);