Implement new GatewaySender, enforceThreadsConnectSameReceiver configuration property.

Resolves gh-480.

Resolves gh-389.
This commit is contained in:
John Blum
2021-01-28 21:21:15 -08:00
committed by John Blum
parent 76bbfcab68
commit a92a78b7f2
4 changed files with 69 additions and 28 deletions

View File

@@ -73,6 +73,7 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
private List<GatewayTransportFilter> 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<Ga
Optional.ofNullable(getDiskSynchronous()).ifPresent(gatewaySenderFactory::setDiskSynchronous);
Optional.ofNullable(getDispatcherThreads()).ifPresent(gatewaySenderFactory::setDispatcherThreads);
Optional.ofNullable(getEnforceThreadsConnectToSameReceiver())
.ifPresent(gatewaySenderFactory::setEnforceThreadsConnectSameReceiver);
CollectionUtils.nullSafeList(getEventFilters()).forEach(gatewaySenderFactory::addGatewayEventFilter);
@@ -259,6 +262,14 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
return this.dispatcherThreads;
}
public void setEnforceThreadsConnectToSameReceiver(Boolean enforceThreadsConnectToSameReceiver) {
this.enforceThreadsConnectToSameReceiver = enforceThreadsConnectToSameReceiver;
}
public Boolean getEnforceThreadsConnectToSameReceiver() {
return this.enforceThreadsConnectToSameReceiver;
}
public void setEventFilters(List<GatewayEventFilter> eventFilters) {
this.eventFilters = eventFilters;
}

View File

@@ -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();
}

View File

@@ -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<GatewayEventFilter> eventFilters;
private List<GatewayTransportFilter> transportFilters;
private final List<GatewayEventFilter> eventFilters = new ArrayList<>();
private final List<GatewayTransportFilter> 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<Boolean>) 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;

View File

@@ -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));
}
}