diff --git a/pom.xml b/pom.xml index 3fb78b5c..235abd55 100644 --- a/pom.xml +++ b/pom.xml @@ -225,7 +225,6 @@ ${spring-shell.version} test - diff --git a/src/main/asciidoc/reference/bootstrap-annotations-quickstart.adoc b/src/main/asciidoc/reference/bootstrap-annotations-quickstart.adoc index 2243f0f9..f7377850 100644 --- a/src/main/asciidoc/reference/bootstrap-annotations-quickstart.adoc +++ b/src/main/asciidoc/reference/bootstrap-annotations-quickstart.adoc @@ -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]. \ No newline at end of file +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]. \ No newline at end of file diff --git a/src/main/asciidoc/reference/bootstrap-annotations.adoc b/src/main/asciidoc/reference/bootstrap-annotations.adoc index 5c3756db..9080fea8 100644 --- a/src/main/asciidoc/reference/bootstrap-annotations.adoc +++ b/src/main/asciidoc/reference/bootstrap-annotations.adoc @@ -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. diff --git a/src/main/java/org/springframework/data/gemfire/config/annotation/EnableGatewaySender.java b/src/main/java/org/springframework/data/gemfire/config/annotation/EnableGatewaySender.java new file mode 100644 index 00000000..2c0ed6fc --- /dev/null +++ b/src/main/java/org/springframework/data/gemfire/config/annotation/EnableGatewaySender.java @@ -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. + *

This property can also be configured using the {@literal spring.data.gemfire.gateway.sender..alert-threshold} + *

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. + *

This property can also be configured using the {@literal spring.data.gemfire.gateway.sender..batch-conflation-enabled} + *

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. + *

This property can also be configured using the {@literal spring.data.gemfire.gateway.sender..batch-size} + *

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}. + *

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. + *

This property can also be configured using the {@literal spring.data.gemfire.gateway.sender..batch-time-interval} + *

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}. + *

This property can also be configured using the {@literal spring.data.gemfire.gateway.sender..diskstore-reference} + *

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. + *

This property can also be configured using the {@literal spring.data.gemfire.gateway.sender..disk-synchronous} + *

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. + *

This property can also be configured using the {@literal spring.data.gemfire.gateway.sender..dispatcher-threads} + *

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}. + *

This property can also be configured using the {@literal spring.data.gemfire.gateway.sender..event-filters} property. + *

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. + *

This property can also be configured using the {@literal spring.data.gemfire.gateway.sender..event-substitution-filter} + *

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 + *

This property can also be configured using the {@literal spring.data.gemfire.gateway.sender..manual-start} + *

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. + *

This property can also be configured using the {@literal spring.data.gemfire.gateway.sender..maximum-queue-memory} + *

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..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}. + *

There are three different ordering policies: + *

+ *

This property can also be configured using the {@literal spring.data.gemfire.gateway.sender..order-policy}

+ *

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}. + *

This property can also be configured using the {@literal spring.data.gemfire.gateway.sender..parallel} + *

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. + *

This property can also be configured using the {@literal spring.data.gemfire.gateway.sender..persistent} + *

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}. + *

This property can also be configured using the {@literal spring.data.gemfire.gateway.sender..region-names} property. + *

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. + *

This property can also be configured using the {@literal spring.data.gemfire.gateway.sender..remote-distributed-system-id} + *

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. + *

This property can also be configured using the {@literal spring.data.gemfire.gateway.sender..socket-buffer-size} property. + *

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). + *

This property can also be configured using the {@literal spring.data.gemfire.gateway.sender..socket-read-timeout} property. + *

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} + *

This property can also be configured using the {@literal spring.data.gemfire.gateway.sender..transport-filters}s property + *

Default value is an empty list + */ + String[] transportFilters() default {}; +} diff --git a/src/main/java/org/springframework/data/gemfire/config/annotation/EnableGatewaySenders.java b/src/main/java/org/springframework/data/gemfire/config/annotation/EnableGatewaySenders.java new file mode 100644 index 00000000..4be7ce8d --- /dev/null +++ b/src/main/java/org/springframework/data/gemfire/config/annotation/EnableGatewaySenders.java @@ -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 gatewaySenders 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. + *

This property can also be configured using the spring.data.gemfire.gateway.sender.alert-threshold + *

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. + *

This property can also be configured using the spring.data.gemfire.gateway.sender.batch-conflation-enabled + *

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 batchSize or batchTimeInterval + * is met. + *

This property can also be configured using the spring.data.gemfire.gateway.sender.batch-size + *

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}. + *

This property works in conjunction with the {@link EnableGatewaySenders#batchSize()} setting. + * The {@link org.apache.geode.cache.wan.GatewaySender} will send when either the batchSize or batchTimeInterval + * is met. + *

This property can also be configured using the spring.data.gemfire.gateway.sender.batch-time-interval + *

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 persistent property is set to true. + *

This property can also be configured using the spring.data.gemfire.gateway.sender.diskstore-reference + *

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. + *

This property can also be configured using the spring.data.gemfire.gateway.sender.disk-synchronous + *

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. + *

This property can also be configured using the spring.data.gemfire.gateway.sender.dispatcher-threads + *

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}. + *

This property can also be configured using the spring.data.gemfire.gateway.sender.event-filters property. + *

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. + *

This property can also be configured using the spring.data.gemfire.gateway.sender.event-substitution-filter + *

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 + *

This property can also be configured using the spring.data.gemfire.gateway.sender.manual-start + *

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. + *

This property can also be configured using the spring.data.gemfire.gateway.sender.maximum-queue-memory + *

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. + *

This property can also be configured using the spring.data.gemfire.gateway.sender.remote-distributed-system-id + *

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}. + *

There are three different ordering policies: + *

+ *

This property can also be configured using the spring.data.gemfire.gateway.sender.order-policy

+ *

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}. + *

This property can also be configured using the spring.data.gemfire.gateway.sender.parallel + *

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 diskStoreReference property. + *

This property can also be configured using the spring.data.gemfire.gateway.sender.persistent + *

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}. + *

This property can also be configured using the spring.data.gemfire.gateway.sender.region-names property. + *

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. + *

This property can also be configured using the spring.data.gemfire.gateway.sender.socket-buffer-size property. + *

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). + *

This property can also be configured using the spring.data.gemfire.gateway.sender.socket-read-timeout property. + *

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} + *

This property can also be configured using the spring.data.gemfire.gateway.sender.transport-filterss property + *

Default value is an empty list + */ + String[] transportFilters() default {}; +} diff --git a/src/main/java/org/springframework/data/gemfire/config/annotation/GatewaySenderConfiguration.java b/src/main/java/org/springframework/data/gemfire/config/annotation/GatewaySenderConfiguration.java new file mode 100644 index 00000000..54d6118e --- /dev/null +++ b/src/main/java/org/springframework/data/gemfire/config/annotation/GatewaySenderConfiguration.java @@ -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 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 annotationAttributesMap = annotationMetadata + .getAnnotationAttributes(EnableGatewaySender.class.getName()); + + AnnotationAttributes gatewaySenderAnnotation = AnnotationAttributes.fromMap(annotationAttributesMap); + + registerGatewaySender(gatewaySenderAnnotation, beanDefinitionRegistry, null); + } + } + + @Override protected Class 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 application.properties. + * + * @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 application.properties. + * + * @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 getValueFromAnnotation(AnnotationAttributes gatewaySenderAnnotation, Supplier supplier) { + if (gatewaySenderAnnotation == null) { + return null; + } + else { + return supplier.get(); + } + } + + /** + * In this method a {@link GatewaySender} is configured from properties defined within an application.properties + * file. + * These properties are "named" properties and will follow the following pattern: + * {@literal spring.data.gemfire.gateway.sender..manual-start} + * + * @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 application.properties + * file. + * These properties are "named" properties and will follow the following pattern: + * {@literal spring.data.gemfire.gateway.sender..manual-start} + * + * @param gatewaySenderBeanBuilder + * @param gatewaySenderName + * @param propertyName + * @param defaultValue + */ + private 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 resolveGatewaySenderConfigurers() { + + return Optional.ofNullable(this.gatewaySenderConfigurers) + .filter(gatewaySenderConfigurers -> !gatewaySenderConfigurers.isEmpty()) + .orElseGet(() -> + Collections.singletonList(LazyResolvingComposableGatewaySenderConfigurer.create(getBeanFactory()))); + } + + private T resolveValueFromProperty(String gatewaySenderName, String propertyName, + T defaultValue) { + Class clazz = (Class) defaultValue.getClass(); + + Optional gatewaySenderProperty = Optional + .ofNullable(resolveProperty(gatewaySenderProperty(propertyName), clazz, null)); + + Optional namedGatewaySenderProperty = Optional.ofNullable( + resolveProperty(namedGatewaySenderProperty(gatewaySenderName, propertyName), clazz, null)); + + return namedGatewaySenderProperty.orElse(gatewaySenderProperty.orElse(null)); + } + + + private BeanDefinitionBuilder setPropertyValueIfNotDefault(BeanDefinitionBuilder beanDefinitionBuilder, + String propertyName, T value, T defaultValue) { + + return beanDefinitionBuilder.addPropertyValue(propertyName, Optional.ofNullable(value).orElse(defaultValue)); + } + + private 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 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) { + + } +} diff --git a/src/main/java/org/springframework/data/gemfire/config/annotation/GatewaySenderConfigurer.java b/src/main/java/org/springframework/data/gemfire/config/annotation/GatewaySenderConfigurer.java new file mode 100644 index 00000000..8b81a4c9 --- /dev/null +++ b/src/main/java/org/springframework/data/gemfire/config/annotation/GatewaySenderConfigurer.java @@ -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 { + + /** + * 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); +} diff --git a/src/main/java/org/springframework/data/gemfire/config/annotation/GatewaySendersConfiguration.java b/src/main/java/org/springframework/data/gemfire/config/annotation/GatewaySendersConfiguration.java new file mode 100644 index 00000000..d0424cee --- /dev/null +++ b/src/main/java/org/springframework/data/gemfire/config/annotation/GatewaySendersConfiguration.java @@ -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 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 getAnnotationType() { + return EnableGatewaySenders.class; + } + + private void registerDefaultGatewaySender(BeanDefinitionRegistry registry, + AnnotationAttributes parentGatewaySendersAnnotation) { + registerGatewaySender("GatewaySender", parentGatewaySendersAnnotation, registry, + parentGatewaySendersAnnotation); + + } +} diff --git a/src/main/java/org/springframework/data/gemfire/config/annotation/LazyResolvingComposableGatewaySenderConfigurer.java b/src/main/java/org/springframework/data/gemfire/config/annotation/LazyResolvingComposableGatewaySenderConfigurer.java new file mode 100644 index 00000000..9e4cbafe --- /dev/null +++ b/src/main/java/org/springframework/data/gemfire/config/annotation/LazyResolvingComposableGatewaySenderConfigurer.java @@ -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 + implements GatewaySenderConfigurer { + + public static LazyResolvingComposableGatewaySenderConfigurer create() { + return create(null); + } + + public static LazyResolvingComposableGatewaySenderConfigurer create(@Nullable BeanFactory beanFactory) { + return new LazyResolvingComposableGatewaySenderConfigurer().with(beanFactory); + } + + @Override + protected Class getConfigurerType() { + return GatewaySenderConfigurer.class; + } +} diff --git a/src/main/java/org/springframework/data/gemfire/config/annotation/support/AbstractAnnotationConfigSupport.java b/src/main/java/org/springframework/data/gemfire/config/annotation/support/AbstractAnnotationConfigSupport.java index a63b04f2..c9e4558e 100644 --- a/src/main/java/org/springframework/data/gemfire/config/annotation/support/AbstractAnnotationConfigSupport.java +++ b/src/main/java/org/springframework/data/gemfire/config/annotation/support/AbstractAnnotationConfigSupport.java @@ -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}. * diff --git a/src/main/java/org/springframework/data/gemfire/config/support/GatewaySenderBeanFactoryPostProcessor.java b/src/main/java/org/springframework/data/gemfire/config/support/GatewaySenderBeanFactoryPostProcessor.java new file mode 100644 index 00000000..c510cd81 --- /dev/null +++ b/src/main/java/org/springframework/data/gemfire/config/support/GatewaySenderBeanFactoryPostProcessor.java @@ -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> cachedBeanDefinitions = populateBeanDefinitionCache(beanFactory); + + //Create a list of gatewaySender to Regions mapping + Map> 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> gatewaySenderToRegions) { + gatewaySenderToRegions.entrySet().forEach(entry -> { + List beanReferenceList = entry.getValue().stream() + .map(RuntimeBeanReference::new).collect(Collectors.toList()); + ManagedList 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> groupGatewaySenderPerRegion( + Map> cachedBeanDefinitions, + Map> gatewaySendersPerRegion) { + + Optional> gatewaySenders = Optional + .ofNullable(cachedBeanDefinitions.get("gatewaySenders")); + + gatewaySenders.ifPresent(gatewaySendersMap -> gatewaySendersMap.forEach((key, gatewaySenderEntryValue) -> { + + PropertyValue regions = gatewaySenderEntryValue.getPropertyValues().getPropertyValue("regions"); + Optional regionsOptional = Optional.ofNullable(regions); + + Collection regionNames; + List 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> populateBeanDefinitionCache( + ConfigurableListableBeanFactory beanFactory) { + Map> 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 <,List> grouped by the `key` + * + * @param beanName + * @param beanDefinition + * @param cachedBeanDefinitions + * @param key + * @return + */ + private Map> addBeanDefinitionToList(String beanName, + BeanDefinition beanDefinition, + Map> cachedBeanDefinitions, String key) { + Map 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> regionMapping, + String gatewaySenderName, Collection regionNames) { + regionNames.forEach(regionName -> { + List 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 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 resolveBeanClass(BeanDefinition beanDefinition, + ConfigurableListableBeanFactory beanFactory) { + + Optional 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; + }); + } +} diff --git a/src/main/java/org/springframework/data/gemfire/wan/AbstractWANComponentFactoryBean.java b/src/main/java/org/springframework/data/gemfire/wan/AbstractWANComponentFactoryBean.java index 612f72fa..d839e450 100644 --- a/src/main/java/org/springframework/data/gemfire/wan/AbstractWANComponentFactoryBean.java +++ b/src/main/java/org/springframework/data/gemfire/wan/AbstractWANComponentFactoryBean.java @@ -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 extends AbstractFactoryBeanSupport implements DisposableBean, InitializingBean { + @Autowired protected Cache cache; protected final Logger logger = LoggerFactory.getLogger(getClass()); @@ -49,12 +50,11 @@ public abstract class AbstractWANComponentFactoryBean 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 extends AbstractFactory this.name = name; } + public void setCache(Cache cache) { + this.cache = cache; + } + public String getName() { return StringUtils.hasText(this.name) diff --git a/src/main/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBean.java b/src/main/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBean.java index ae44c9bb..f38293a4 100644 --- a/src/main/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBean.java +++ b/src/main/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBean.java @@ -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 gatewaySenderConfigurers = Collections.emptyList(); + + private List 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 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 getEventFilters() { + return eventFilters; + } + + public void setEventFilters(List eventFilters) { + this.eventFilters = eventFilters; + } + + public List getTransportFilters() { + return transportFilters; + } + + public void setTransportFilters(List 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 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 gatewayTransportFilters) { - this.transportFilters = gatewayTransportFilters; + public String getDiskStoreReference() { + return diskStoreReference; + } + + public void setDiskStoreReference(String diskStoreReference) { + this.diskStoreReference = diskStoreReference; + } + + public void setGatewaySenderConfigurers(List gatewaySenderConfigurers) { + this.gatewaySenderConfigurers = gatewaySenderConfigurers; + } + + public void setDiskStoreRef(String diskStoreRef) { + this.diskStoreReference = diskStoreRef; + } + + private List getRegions() { + return regions; + } + + public void setRegions(List regions) { + this.regions = regions; + } + + public void setRegions(String[] regions) { + this.regions = Arrays.asList(regions); } } diff --git a/src/main/java/org/springframework/data/gemfire/wan/OrderPolicyConverter.java b/src/main/java/org/springframework/data/gemfire/wan/OrderPolicyConverter.java index 2a571fbf..271ce846 100644 --- a/src/main/java/org/springframework/data/gemfire/wan/OrderPolicyConverter.java +++ b/src/main/java/org/springframework/data/gemfire/wan/OrderPolicyConverter.java @@ -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 { +public class OrderPolicyConverter extends AbstractPropertyEditorConverterSupport { /** * 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); + } } diff --git a/src/main/java/org/springframework/data/gemfire/wan/OrderPolicyType.java b/src/main/java/org/springframework/data/gemfire/wan/OrderPolicyType.java index 30b50461..a66581e5 100644 --- a/src/main/java/org/springframework/data/gemfire/wan/OrderPolicyType.java +++ b/src/main/java/org/springframework/data/gemfire/wan/OrderPolicyType.java @@ -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; } } diff --git a/src/test/java/org/springframework/data/gemfire/config/annotation/GatewaySenderConfigurationTests.java b/src/test/java/org/springframework/data/gemfire/config/annotation/GatewaySenderConfigurationTests.java new file mode 100644 index 00000000..4ff98782 --- /dev/null +++ b/src/test/java/org/springframework/data/gemfire/config/annotation/GatewaySenderConfigurationTests.java @@ -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 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 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 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 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; + } + } +} diff --git a/src/test/java/org/springframework/data/gemfire/config/annotation/GatewaySenderConfigurerTests.java b/src/test/java/org/springframework/data/gemfire/config/annotation/GatewaySenderConfigurerTests.java new file mode 100644 index 00000000..f4479f38 --- /dev/null +++ b/src/test/java/org/springframework/data/gemfire/config/annotation/GatewaySenderConfigurerTests.java @@ -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 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); + }); + } + } +} diff --git a/src/test/java/org/springframework/data/gemfire/config/annotation/GatewaySenderPropertiesTests.java b/src/test/java/org/springframework/data/gemfire/config/annotation/GatewaySenderPropertiesTests.java new file mode 100644 index 00000000..e8771649 --- /dev/null +++ b/src/test/java/org/springframework/data/gemfire/config/annotation/GatewaySenderPropertiesTests.java @@ -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 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; + } + } +} diff --git a/src/test/java/org/springframework/data/gemfire/repository/sample/Customer.java b/src/test/java/org/springframework/data/gemfire/repository/sample/Customer.java index 4b9c6796..b6500af1 100644 --- a/src/test/java/org/springframework/data/gemfire/repository/sample/Customer.java +++ b/src/test/java/org/springframework/data/gemfire/repository/sample/Customer.java @@ -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 diff --git a/src/test/java/org/springframework/data/gemfire/test/StubGatewaySenderFactory.java b/src/test/java/org/springframework/data/gemfire/test/StubGatewaySenderFactory.java index ecf9d467..8a273e2d 100644 --- a/src/test/java/org/springframework/data/gemfire/test/StubGatewaySenderFactory.java +++ b/src/test/java/org/springframework/data/gemfire/test/StubGatewaySenderFactory.java @@ -111,11 +111,11 @@ public class StubGatewaySenderFactory implements GatewaySenderFactory { } }); doAnswer(new Answer() { - 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; } diff --git a/src/test/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBeanTest.java b/src/test/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBeanTest.java index e363c0f5..51322acc 100644 --- a/src/test/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBeanTest.java +++ b/src/test/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBeanTest.java @@ -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);