From a92a78b7f295bae95e36950c2f1e3b89a112f6d6 Mon Sep 17 00:00:00 2001 From: John Blum Date: Thu, 28 Jan 2021 21:21:15 -0800 Subject: [PATCH] Implement new GatewaySender, enforceThreadsConnectSameReceiver configuration property. Resolves gh-480. Resolves gh-389. --- .../gemfire/wan/GatewaySenderFactoryBean.java | 11 ++++ .../gemfire/wan/GatewaySenderWrapper.java | 9 +++ .../test/StubGatewaySenderFactory.java | 17 +++--- .../wan/GatewaySenderFactoryBeanTest.java | 60 ++++++++++++------- 4 files changed, 69 insertions(+), 28 deletions(-) diff --git a/spring-data-geode/src/main/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBean.java b/spring-data-geode/src/main/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBean.java index fd02bc8c..87813647 100644 --- a/spring-data-geode/src/main/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBean.java +++ b/spring-data-geode/src/main/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBean.java @@ -73,6 +73,7 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean transportFilters; private Boolean diskSynchronous; + private Boolean enforceThreadsConnectToSameReceiver; private Boolean batchConflationEnabled; private Boolean groupTransactionEvents; private Boolean parallel; @@ -133,6 +134,8 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean eventFilters) { this.eventFilters = eventFilters; } diff --git a/spring-data-geode/src/main/java/org/springframework/data/gemfire/wan/GatewaySenderWrapper.java b/spring-data-geode/src/main/java/org/springframework/data/gemfire/wan/GatewaySenderWrapper.java index 9d8bdebb..cf3e0342 100644 --- a/spring-data-geode/src/main/java/org/springframework/data/gemfire/wan/GatewaySenderWrapper.java +++ b/spring-data-geode/src/main/java/org/springframework/data/gemfire/wan/GatewaySenderWrapper.java @@ -143,6 +143,14 @@ public class GatewaySenderWrapper implements GatewaySender { return delegate.getDispatcherThreads(); } + /** + * @inheritDoc + */ + @Override + public boolean getEnforceThreadsConnectSameReceiver() { + return this.delegate.getEnforceThreadsConnectSameReceiver(); + } + /** * @inheritDoc */ @@ -155,6 +163,7 @@ public class GatewaySenderWrapper implements GatewaySender { * @inheritDoc */ @Override + @SuppressWarnings("rawtypes") public GatewayEventSubstitutionFilter getGatewayEventSubstitutionFilter() { return this.delegate.getGatewayEventSubstitutionFilter(); } diff --git a/spring-data-geode/src/test/java/org/springframework/data/gemfire/test/StubGatewaySenderFactory.java b/spring-data-geode/src/test/java/org/springframework/data/gemfire/test/StubGatewaySenderFactory.java index 3edb6282..8171dcd4 100644 --- a/spring-data-geode/src/test/java/org/springframework/data/gemfire/test/StubGatewaySenderFactory.java +++ b/spring-data-geode/src/test/java/org/springframework/data/gemfire/test/StubGatewaySenderFactory.java @@ -45,6 +45,7 @@ public class StubGatewaySenderFactory implements GatewaySenderFactory { private boolean parallel; private boolean persistenceEnabled; private boolean running = false; + private boolean enforceThreadsConnectSameReceiver; private int alertThreshold; private int batchSize; @@ -57,18 +58,13 @@ public class StubGatewaySenderFactory implements GatewaySenderFactory { private GatewayEventSubstitutionFilter gatewayEventSubstitutionFilter; - private List eventFilters; - private List transportFilters; + private final List eventFilters = new ArrayList<>(); + private final List transportFilters = new ArrayList<>(); private OrderPolicy orderPolicy; private String diskStoreName; - public StubGatewaySenderFactory() { - this.eventFilters = new ArrayList<>(); - this.transportFilters = new ArrayList<>(); - } - @Override public GatewaySenderFactory addGatewayEventFilter(GatewayEventFilter filter) { eventFilters.add(filter); @@ -105,7 +101,6 @@ public class StubGatewaySenderFactory implements GatewaySenderFactory { when(gatewaySender.isParallel()).thenReturn(this.parallel); when(gatewaySender.isPersistenceEnabled()).thenReturn(this.persistenceEnabled); when(gatewaySender.getOrderPolicy()).thenReturn(this.orderPolicy); - when(gatewaySender.mustGroupTransactionEvents()).thenReturn(this.groupTransactionEvents); when(gatewaySender.isRunning()).thenAnswer((Answer) invocation -> this.running); doAnswer(invocation -> { running = true; @@ -163,6 +158,12 @@ public class StubGatewaySenderFactory implements GatewaySenderFactory { return this; } + @Override + public GatewaySenderFactory setEnforceThreadsConnectSameReceiver(boolean b) { + this.enforceThreadsConnectSameReceiver = b; + return this; + } + @Override public GatewaySenderFactory setGroupTransactionEvents(boolean groupTransactionEvents) { this.groupTransactionEvents = groupTransactionEvents; diff --git a/spring-data-geode/src/test/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBeanTest.java b/spring-data-geode/src/test/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBeanTest.java index 2cc9b08f..9a2cc82a 100644 --- a/spring-data-geode/src/test/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBeanTest.java +++ b/spring-data-geode/src/test/java/org/springframework/data/gemfire/wan/GatewaySenderFactoryBeanTest.java @@ -15,11 +15,11 @@ */ package org.springframework.data.gemfire.wan; -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertNotNull; +import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; +import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @@ -128,9 +128,9 @@ public class GatewaySenderFactoryBeanTest { GatewaySender gatewaySender = factoryBean.getObject(); - assertNotNull(gatewaySender); - assertEquals("g0", gatewaySender.getId()); - assertEquals(69, gatewaySender.getRemoteDSId()); + assertThat(gatewaySender).isNotNull(); + assertThat(gatewaySender.getId()).isEqualTo("g0"); + assertThat(gatewaySender.getRemoteDSId()).isEqualTo(69); } @Test @@ -152,9 +152,9 @@ public class GatewaySenderFactoryBeanTest { GatewaySender gatewaySender = factoryBean.getObject(); - assertNotNull(gatewaySender); - assertEquals("g4", gatewaySender.getId()); - assertEquals(21, gatewaySender.getRemoteDSId()); + assertThat(gatewaySender).isNotNull(); + assertThat(gatewaySender.getId()).isEqualTo("g4"); + assertThat(gatewaySender.getRemoteDSId()).isEqualTo(21); } @Test @@ -175,9 +175,9 @@ public class GatewaySenderFactoryBeanTest { GatewaySender gatewaySender = factoryBean.getObject(); - assertNotNull(gatewaySender); - assertEquals("g1", gatewaySender.getId()); - assertEquals(69, gatewaySender.getRemoteDSId()); + assertThat(gatewaySender).isNotNull(); + assertThat(gatewaySender.getId()).isEqualTo("g1"); + assertThat(gatewaySender.getRemoteDSId()).isEqualTo(69); } @Test @@ -198,9 +198,9 @@ public class GatewaySenderFactoryBeanTest { GatewaySender gatewaySender = factoryBean.getObject(); - assertNotNull(gatewaySender); - assertEquals("g7", gatewaySender.getId()); - assertEquals(51, gatewaySender.getRemoteDSId()); + assertThat(gatewaySender).isNotNull(); + assertThat(gatewaySender.getId()).isEqualTo("g7"); + assertThat(gatewaySender.getRemoteDSId()).isEqualTo(51); } @Test @@ -222,9 +222,9 @@ public class GatewaySenderFactoryBeanTest { GatewaySender gatewaySender = factoryBean.getObject(); - assertNotNull(gatewaySender); - assertEquals("g6", gatewaySender.getId()); - assertEquals(51, gatewaySender.getRemoteDSId()); + assertThat(gatewaySender).isNotNull(); + assertThat(gatewaySender.getId()).isEqualTo("g6"); + assertThat(gatewaySender.getRemoteDSId()).isEqualTo(51); } @Test @@ -246,8 +246,28 @@ public class GatewaySenderFactoryBeanTest { GatewaySender gatewaySender = factoryBean.getObject(); - assertNotNull(gatewaySender); - assertEquals("g5", gatewaySender.getId()); - assertEquals(42, gatewaySender.getRemoteDSId()); + assertThat(gatewaySender).isNotNull(); + assertThat(gatewaySender.getId()).isEqualTo("g5"); + assertThat(gatewaySender.getRemoteDSId()).isEqualTo(42); + } + + @Test + public void gatewaySenderFactoryBeanSetsGatewaySenderFactoryEnforceThreadsConnectSameReceiver() { + + GatewaySenderFactory mockGatewaySenderFactory = + mockGatewaySenderFactory("g10", 69); + + GatewaySenderFactoryBean factoryBean = + new GatewaySenderFactoryBean(mockCacheWithGatewayInfrastructure(mockGatewaySenderFactory)); + + factoryBean.setName("g10"); + factoryBean.setRemoteDistributedSystemId(69); + factoryBean.setEnforceThreadsConnectToSameReceiver(true); + factoryBean.doInit(); + + assertThat(factoryBean.getEnforceThreadsConnectToSameReceiver()).isTrue(); + + verify(mockGatewaySenderFactory, times(1)) + .setEnforceThreadsConnectSameReceiver(eq(true)); } }