SGF-199 fixed GatewaySender manual-start
This commit is contained in:
2
.gitignore
vendored
2
.gitignore
vendored
@@ -12,6 +12,8 @@ pom.xml
|
|||||||
.classpath
|
.classpath
|
||||||
.project
|
.project
|
||||||
.settings/
|
.settings/
|
||||||
|
.idea/
|
||||||
|
_site/
|
||||||
|
|
||||||
/samples/hello-world/vf.gf.dmn-events.cfg
|
/samples/hello-world/vf.gf.dmn-events.cfg
|
||||||
/samples/hello-world/vf.gf.dmn-license.cfg
|
/samples/hello-world/vf.gf.dmn-license.cfg
|
||||||
|
|||||||
@@ -24,6 +24,7 @@ import org.springframework.beans.factory.DisposableBean;
|
|||||||
import org.springframework.context.SmartLifecycle;
|
import org.springframework.context.SmartLifecycle;
|
||||||
import org.springframework.core.io.Resource;
|
import org.springframework.core.io.Resource;
|
||||||
import org.springframework.data.gemfire.client.ClientRegionFactoryBean;
|
import org.springframework.data.gemfire.client.ClientRegionFactoryBean;
|
||||||
|
import org.springframework.data.gemfire.wan.GatewaySenderWrapper;
|
||||||
import org.springframework.util.Assert;
|
import org.springframework.util.Assert;
|
||||||
import org.springframework.util.ObjectUtils;
|
import org.springframework.util.ObjectUtils;
|
||||||
import org.springframework.util.ReflectionUtils;
|
import org.springframework.util.ReflectionUtils;
|
||||||
@@ -411,7 +412,9 @@ public class RegionFactoryBean<K, V> extends RegionLookupFactoryBean<K, V> imple
|
|||||||
synchronized (gatewaySenders) {
|
synchronized (gatewaySenders) {
|
||||||
for (Object obj : gatewaySenders) {
|
for (Object obj : gatewaySenders) {
|
||||||
GatewaySender gws = (GatewaySender) obj;
|
GatewaySender gws = (GatewaySender) obj;
|
||||||
|
System.out.println("manual start = " + gws.isManualStart());
|
||||||
if (!gws.isManualStart() && !gws.isRunning()) {
|
if (!gws.isManualStart() && !gws.isRunning()) {
|
||||||
|
System.out.println("starting ...");
|
||||||
gws.start();
|
gws.start();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -18,6 +18,7 @@ package org.springframework.data.gemfire.wan;
|
|||||||
import java.util.Arrays;
|
import java.util.Arrays;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
|
|
||||||
|
import org.springframework.beans.DirectFieldAccessor;
|
||||||
import org.springframework.context.SmartLifecycle;
|
import org.springframework.context.SmartLifecycle;
|
||||||
import org.springframework.util.Assert;
|
import org.springframework.util.Assert;
|
||||||
import org.springframework.util.CollectionUtils;
|
import org.springframework.util.CollectionUtils;
|
||||||
@@ -165,7 +166,9 @@ implements SmartLifecycle {
|
|||||||
if (socketReadTimeout != null) {
|
if (socketReadTimeout != null) {
|
||||||
gatewaySenderFactory.setSocketReadTimeout(socketReadTimeout);
|
gatewaySenderFactory.setSocketReadTimeout(socketReadTimeout);
|
||||||
}
|
}
|
||||||
gatewaySender = gatewaySenderFactory.create(getName(), remoteDistributedSystemId);
|
GatewaySenderWrapper wrapper = new GatewaySenderWrapper(gatewaySenderFactory.create(getName(), remoteDistributedSystemId));
|
||||||
|
wrapper.setManualStart(manualStart);
|
||||||
|
gatewaySender = wrapper;
|
||||||
}
|
}
|
||||||
|
|
||||||
public void setRemoteDistributedSystemId(int remoteDistributedSystemId) {
|
public void setRemoteDistributedSystemId(int remoteDistributedSystemId) {
|
||||||
|
|||||||
@@ -0,0 +1,154 @@
|
|||||||
|
package org.springframework.data.gemfire.wan;
|
||||||
|
|
||||||
|
import com.gemstone.gemfire.cache.util.Gateway;
|
||||||
|
import com.gemstone.gemfire.cache.wan.GatewayEventFilter;
|
||||||
|
import com.gemstone.gemfire.cache.wan.GatewaySender;
|
||||||
|
import com.gemstone.gemfire.cache.wan.GatewayTransportFilter;
|
||||||
|
|
||||||
|
import java.util.List;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Created by dturanski on 9/16/13.
|
||||||
|
*/
|
||||||
|
public class GatewaySenderWrapper implements GatewaySender {
|
||||||
|
private final GatewaySender delegate;
|
||||||
|
private boolean manualStart;
|
||||||
|
|
||||||
|
public GatewaySenderWrapper(GatewaySender delegate) {
|
||||||
|
this.delegate = delegate;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void start() {
|
||||||
|
delegate.start();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void stop() {
|
||||||
|
delegate.stop();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void pause() {
|
||||||
|
delegate.pause();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void resume() {
|
||||||
|
delegate.resume();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean isRunning() {
|
||||||
|
return delegate.isRunning();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean isPaused() {
|
||||||
|
return delegate.isPaused();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void addGatewayEventFilter(GatewayEventFilter filter) {
|
||||||
|
delegate.addGatewayEventFilter(filter);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void removeGatewayEventFilter(GatewayEventFilter filter) {
|
||||||
|
delegate.removeGatewayEventFilter(filter);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public String getId() {
|
||||||
|
return delegate.getId();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public int getRemoteDSId() {
|
||||||
|
return delegate.getRemoteDSId();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public int getSocketBufferSize() {
|
||||||
|
return delegate.getSocketBufferSize();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public int getSocketReadTimeout() {
|
||||||
|
return delegate.getSocketReadTimeout();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public String getDiskStoreName() {
|
||||||
|
return delegate.getDiskStoreName();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public int getMaximumQueueMemory() {
|
||||||
|
return delegate.getMaximumQueueMemory();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public int getBatchSize() {
|
||||||
|
return delegate.getBatchSize();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public int getBatchTimeInterval() {
|
||||||
|
return delegate.getBatchTimeInterval();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean isBatchConflationEnabled() {
|
||||||
|
return delegate.isBatchConflationEnabled();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean isPersistenceEnabled() {
|
||||||
|
return delegate.isPersistenceEnabled();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public int getAlertThreshold() {
|
||||||
|
return delegate.getAlertThreshold();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public List<GatewayEventFilter> getGatewayEventFilters() {
|
||||||
|
return delegate.getGatewayEventFilters();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public List<GatewayTransportFilter> getGatewayTransportFilters() {
|
||||||
|
return delegate.getGatewayTransportFilters();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean isDiskSynchronous() {
|
||||||
|
return delegate.isDiskSynchronous();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean isManualStart() {
|
||||||
|
return this.manualStart;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean isParallel() {
|
||||||
|
return delegate.isParallel();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public int getDispatcherThreads() {
|
||||||
|
return delegate.getDispatcherThreads();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public Gateway.OrderPolicy getOrderPolicy() {
|
||||||
|
return delegate.getOrderPolicy();
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setManualStart(boolean manualStart) {
|
||||||
|
this.manualStart = manualStart;
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -148,7 +148,8 @@ public class GemfireV7GatewayNamespaceTest extends RecreatingContextTest {
|
|||||||
assertTrue(transportFilters.get(0) instanceof TestTransportFilter);
|
assertTrue(transportFilters.get(0) instanceof TestTransportFilter);
|
||||||
|
|
||||||
assertEquals(1, gws.getRemoteDSId());
|
assertEquals(1, gws.getRemoteDSId());
|
||||||
assertEquals(true, gws.isManualStart());
|
assertEquals(false, gws.isManualStart());
|
||||||
|
assertEquals(true,gws.isRunning());
|
||||||
assertEquals(10, gws.getAlertThreshold());
|
assertEquals(10, gws.getAlertThreshold());
|
||||||
assertEquals(11, gws.getBatchSize());
|
assertEquals(11, gws.getBatchSize());
|
||||||
assertEquals(3000, gws.getBatchTimeInterval());
|
assertEquals(3000, gws.getBatchTimeInterval());
|
||||||
|
|||||||
@@ -12,17 +12,23 @@
|
|||||||
*/
|
*/
|
||||||
package org.springframework.data.gemfire.test;
|
package org.springframework.data.gemfire.test;
|
||||||
|
|
||||||
|
import static org.mockito.Matchers.anyString;
|
||||||
|
import static org.mockito.Mockito.doAnswer;
|
||||||
import static org.mockito.Mockito.mock;
|
import static org.mockito.Mockito.mock;
|
||||||
import static org.mockito.Mockito.when;
|
import static org.mockito.Mockito.when;
|
||||||
|
|
||||||
import java.util.ArrayList;
|
import java.util.ArrayList;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
|
|
||||||
|
import com.gemstone.gemfire.cache.Region;
|
||||||
import com.gemstone.gemfire.cache.util.Gateway.OrderPolicy;
|
import com.gemstone.gemfire.cache.util.Gateway.OrderPolicy;
|
||||||
import com.gemstone.gemfire.cache.wan.GatewayEventFilter;
|
import com.gemstone.gemfire.cache.wan.GatewayEventFilter;
|
||||||
import com.gemstone.gemfire.cache.wan.GatewaySender;
|
import com.gemstone.gemfire.cache.wan.GatewaySender;
|
||||||
import com.gemstone.gemfire.cache.wan.GatewaySenderFactory;
|
import com.gemstone.gemfire.cache.wan.GatewaySenderFactory;
|
||||||
import com.gemstone.gemfire.cache.wan.GatewayTransportFilter;
|
import com.gemstone.gemfire.cache.wan.GatewayTransportFilter;
|
||||||
|
import com.sun.org.apache.xpath.internal.operations.Bool;
|
||||||
|
import org.mockito.invocation.InvocationOnMock;
|
||||||
|
import org.mockito.stubbing.Answer;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* @author David Turanski
|
* @author David Turanski
|
||||||
@@ -49,6 +55,8 @@ public class StubGatewaySenderFactory implements GatewaySenderFactory {
|
|||||||
|
|
||||||
private String name;
|
private String name;
|
||||||
|
|
||||||
|
private boolean running = true;
|
||||||
|
|
||||||
private int remoteSystemId;
|
private int remoteSystemId;
|
||||||
|
|
||||||
public StubGatewaySenderFactory() {
|
public StubGatewaySenderFactory() {
|
||||||
@@ -92,6 +100,18 @@ public class StubGatewaySenderFactory implements GatewaySenderFactory {
|
|||||||
when(gatewaySender.isDiskSynchronous()).thenReturn(this.diskSynchronous);
|
when(gatewaySender.isDiskSynchronous()).thenReturn(this.diskSynchronous);
|
||||||
when(gatewaySender.isParallel()).thenReturn(this.parallel);
|
when(gatewaySender.isParallel()).thenReturn(this.parallel);
|
||||||
when(gatewaySender.isPersistenceEnabled()).thenReturn(this.persistenceEnabled);
|
when(gatewaySender.isPersistenceEnabled()).thenReturn(this.persistenceEnabled);
|
||||||
|
doAnswer(new Answer<Void>() {
|
||||||
|
public Void answer(InvocationOnMock invocation) {
|
||||||
|
running = true;
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
}).when(gatewaySender).start();
|
||||||
|
when(gatewaySender.isRunning()).thenAnswer(new Answer<Boolean>(){
|
||||||
|
@Override
|
||||||
|
public Boolean answer(InvocationOnMock invocation) throws Throwable {
|
||||||
|
return running;
|
||||||
|
}
|
||||||
|
});
|
||||||
return gatewaySender;
|
return gatewaySender;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -12,7 +12,7 @@
|
|||||||
|
|
||||||
<gfe:partitioned-region id="region-inner-gateway-sender" >
|
<gfe:partitioned-region id="region-inner-gateway-sender" >
|
||||||
<gfe:gateway-sender
|
<gfe:gateway-sender
|
||||||
manual-start="true"
|
manual-start="false"
|
||||||
remote-distributed-system-id="1"
|
remote-distributed-system-id="1"
|
||||||
alert-threshold="10"
|
alert-threshold="10"
|
||||||
batch-size="11"
|
batch-size="11"
|
||||||
|
|||||||
Reference in New Issue
Block a user