SGF-570 - Respect manual-start on Gateway Senders/Receivers but no longer couple the start/stop lifecycle to the Spring container.
This commit is contained in:
@@ -15,22 +15,20 @@
|
||||
*/
|
||||
package org.springframework.data.gemfire.wan;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.util.List;
|
||||
|
||||
import org.springframework.context.SmartLifecycle;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.CollectionUtils;
|
||||
import org.springframework.util.StringUtils;
|
||||
|
||||
import com.gemstone.gemfire.cache.Cache;
|
||||
import com.gemstone.gemfire.cache.wan.GatewayReceiver;
|
||||
import com.gemstone.gemfire.cache.wan.GatewayReceiverFactory;
|
||||
import com.gemstone.gemfire.cache.wan.GatewayTransportFilter;
|
||||
|
||||
import org.springframework.data.gemfire.util.CollectionUtils;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.StringUtils;
|
||||
|
||||
/**
|
||||
* Spring FactoryBean for creating a GemFire {@link GatewayReceiver}.
|
||||
*
|
||||
*
|
||||
* @author David Turanski
|
||||
* @author John Blum
|
||||
* @see org.springframework.context.SmartLifecycle
|
||||
@@ -41,8 +39,7 @@ import com.gemstone.gemfire.cache.wan.GatewayTransportFilter;
|
||||
* @since 1.2.2
|
||||
*/
|
||||
@SuppressWarnings("unused")
|
||||
public class GatewayReceiverFactoryBean extends AbstractWANComponentFactoryBean<GatewayReceiver>
|
||||
implements SmartLifecycle {
|
||||
public class GatewayReceiverFactoryBean extends AbstractWANComponentFactoryBean<GatewayReceiver> {
|
||||
|
||||
private boolean manualStart = false;
|
||||
|
||||
@@ -59,35 +56,39 @@ public class GatewayReceiverFactoryBean extends AbstractWANComponentFactoryBean<
|
||||
private String hostnameForSenders;
|
||||
|
||||
/**
|
||||
* Constructs an instance of the GatewayReceiverFactoryBean class for configuring an initializing
|
||||
* a GemFire Gateway Receiver.
|
||||
* Constructs an instance of the {@link GatewayReceiverFactoryBean} class for configuring and initializing
|
||||
* a GemFire {@link GatewayReceiver}.
|
||||
*
|
||||
* @param cache a reference to the GemFire Cache used to setup the Gateway Receiver.
|
||||
* @param cache reference to the GemFire {@link Cache} used to create the {@link GatewayReceiver}.
|
||||
* @see com.gemstone.gemfire.cache.Cache
|
||||
*/
|
||||
public GatewayReceiverFactoryBean(Cache cache) {
|
||||
super(cache);
|
||||
}
|
||||
|
||||
/**
|
||||
* @inheritDoc
|
||||
*/
|
||||
@Override
|
||||
public GatewayReceiver getObject() throws Exception {
|
||||
return gatewayReceiver;
|
||||
}
|
||||
|
||||
/**
|
||||
* @inheritDoc
|
||||
*/
|
||||
@Override
|
||||
public Class<?> getObjectType() {
|
||||
return (gatewayReceiver != null ? gatewayReceiver.getClass() : GatewayReceiver.class);
|
||||
}
|
||||
|
||||
/**
|
||||
* @inheritDoc
|
||||
*/
|
||||
@Override
|
||||
protected void doInit() throws Exception {
|
||||
GatewayReceiverFactory gatewayReceiverFactory = cache.createGatewayReceiverFactory();
|
||||
|
||||
if (!CollectionUtils.isEmpty(transportFilters)) {
|
||||
for (GatewayTransportFilter transportFilter : transportFilters) {
|
||||
gatewayReceiverFactory.addGatewayTransportFilter(transportFilter);
|
||||
}
|
||||
}
|
||||
|
||||
if (StringUtils.hasText(bindAddress)) {
|
||||
gatewayReceiverFactory.setBindAddress(bindAddress);
|
||||
}
|
||||
@@ -96,28 +97,36 @@ public class GatewayReceiverFactoryBean extends AbstractWANComponentFactoryBean<
|
||||
gatewayReceiverFactory.setHostnameForSenders(hostnameForSenders);
|
||||
}
|
||||
|
||||
if (maximumTimeBetweenPings != null) {
|
||||
gatewayReceiverFactory.setMaximumTimeBetweenPings(maximumTimeBetweenPings);
|
||||
}
|
||||
int localStartPort = defaultPort(startPort, GatewayReceiver.DEFAULT_START_PORT);
|
||||
int localEndPort = defaultPort(endPort, GatewayReceiver.DEFAULT_END_PORT);
|
||||
|
||||
int localStartPort = (startPort != null ? startPort : GatewayReceiver.DEFAULT_START_PORT);
|
||||
int localEndPort = (endPort != null ? endPort : GatewayReceiver.DEFAULT_END_PORT);
|
||||
|
||||
Assert.isTrue(localStartPort <= localEndPort, String.format("'startPort' must be less than or equal to %1$d.",
|
||||
localEndPort));
|
||||
Assert.isTrue(localStartPort <= localEndPort,
|
||||
String.format("'startPort' must be less than or equal to %d.", localEndPort));
|
||||
|
||||
gatewayReceiverFactory.setStartPort(localStartPort);
|
||||
gatewayReceiverFactory.setEndPort(localEndPort);
|
||||
gatewayReceiverFactory.setManualStart(true);
|
||||
gatewayReceiverFactory.setManualStart(manualStart);
|
||||
|
||||
if (maximumTimeBetweenPings != null) {
|
||||
gatewayReceiverFactory.setMaximumTimeBetweenPings(maximumTimeBetweenPings);
|
||||
}
|
||||
|
||||
if (socketBufferSize != null) {
|
||||
gatewayReceiverFactory.setSocketBufferSize(socketBufferSize);
|
||||
}
|
||||
|
||||
for (GatewayTransportFilter transportFilter : CollectionUtils.nullSafeList(transportFilters)) {
|
||||
gatewayReceiverFactory.addGatewayTransportFilter(transportFilter);
|
||||
}
|
||||
|
||||
gatewayReceiver = gatewayReceiverFactory.create();
|
||||
}
|
||||
|
||||
public void setGatewayReceiver(final GatewayReceiver gatewayReceiver) {
|
||||
protected int defaultPort(Integer port, int defaultPort) {
|
||||
return (port != null ? port : defaultPort);
|
||||
}
|
||||
|
||||
public void setGatewayReceiver(GatewayReceiver gatewayReceiver) {
|
||||
this.gatewayReceiver = gatewayReceiver;
|
||||
}
|
||||
|
||||
@@ -152,45 +161,4 @@ public class GatewayReceiverFactoryBean extends AbstractWANComponentFactoryBean<
|
||||
public void setTransportFilters(List<GatewayTransportFilter> transportFilters) {
|
||||
this.transportFilters = transportFilters;
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean isAutoStartup() {
|
||||
return !manualStart;
|
||||
}
|
||||
|
||||
@Override
|
||||
public int getPhase() {
|
||||
return Integer.MAX_VALUE;
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean isRunning() {
|
||||
return gatewayReceiver.isRunning();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void start() {
|
||||
Assert.state(gatewayReceiver != null, "The GatewayReceiver was not properly configured and initialized!");
|
||||
|
||||
if (!isRunning()) {
|
||||
try {
|
||||
gatewayReceiver.start();
|
||||
}
|
||||
catch (IOException e) {
|
||||
throw new RuntimeException("Failed to start Gateway Receiver due to I/O error.", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void stop() {
|
||||
gatewayReceiver.stop();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void stop(final Runnable callback) {
|
||||
stop();
|
||||
callback.run();
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -17,10 +17,6 @@ package org.springframework.data.gemfire.wan;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
import org.springframework.context.SmartLifecycle;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.CollectionUtils;
|
||||
|
||||
import com.gemstone.gemfire.cache.Cache;
|
||||
import com.gemstone.gemfire.cache.util.Gateway;
|
||||
import com.gemstone.gemfire.cache.wan.GatewayEventFilter;
|
||||
@@ -29,22 +25,23 @@ import com.gemstone.gemfire.cache.wan.GatewaySender;
|
||||
import com.gemstone.gemfire.cache.wan.GatewaySenderFactory;
|
||||
import com.gemstone.gemfire.cache.wan.GatewayTransportFilter;
|
||||
|
||||
import org.springframework.beans.factory.FactoryBean;
|
||||
import org.springframework.data.gemfire.util.CollectionUtils;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
/**
|
||||
* FactoryBean for creating a parallel or serial GemFire {@link GatewaySender}.
|
||||
* Spring {@link FactoryBean} for creating a parallel or serial GemFire {@link GatewaySender}.
|
||||
*
|
||||
* @author David Turanski
|
||||
* @author John Blum
|
||||
* @see org.springframework.context.SmartLifecycle
|
||||
* @see org.springframework.data.gemfire.wan.AbstractWANComponentFactoryBean
|
||||
* @see com.gemstone.gemfire.cache.Cache
|
||||
* @see com.gemstone.gemfire.cache.util.Gateway
|
||||
* @see com.gemstone.gemfire.cache.wan.GatewaySender
|
||||
* @see com.gemstone.gemfire.cache.wan.GatewaySenderFactory
|
||||
* @since 1.2.2
|
||||
*/
|
||||
@SuppressWarnings("unused")
|
||||
public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<GatewaySender>
|
||||
implements SmartLifecycle {
|
||||
public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<GatewaySender> {
|
||||
|
||||
private boolean manualStart = false;
|
||||
|
||||
@@ -75,25 +72,35 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
|
||||
private String orderPolicy;
|
||||
|
||||
/**
|
||||
* Constructs an instance of the GatewaySenderFactoryBean class initialized with a reference to the GemFire cache.
|
||||
* Constructs an instance of the {@link GatewaySenderFactoryBean} class initialized with
|
||||
* a reference to the GemFire {@link Cache}.
|
||||
*
|
||||
* @param cache the Gemfire cache reference.
|
||||
* @param cache reference to the GemFire {@link Cache} used to create the GemFire {@link GatewaySender}.
|
||||
* @see com.gemstone.gemfire.cache.Cache
|
||||
*/
|
||||
public GatewaySenderFactoryBean(final Cache cache) {
|
||||
public GatewaySenderFactoryBean(Cache cache) {
|
||||
super(cache);
|
||||
}
|
||||
|
||||
/**
|
||||
* @inheritDoc
|
||||
*/
|
||||
@Override
|
||||
public GatewaySender getObject() throws Exception {
|
||||
return gatewaySender;
|
||||
}
|
||||
|
||||
/**
|
||||
* @inheritDoc
|
||||
*/
|
||||
@Override
|
||||
public Class<?> getObjectType() {
|
||||
return (gatewaySender != null ? gatewaySender.getClass() : GatewaySender.class);
|
||||
}
|
||||
|
||||
/**
|
||||
* @inheritDoc
|
||||
*/
|
||||
@Override
|
||||
protected void doInit() {
|
||||
GatewaySenderFactory gatewaySenderFactory = (this.factory != null ? (GatewaySenderFactory) factory
|
||||
@@ -103,6 +110,10 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
|
||||
gatewaySenderFactory.setAlertThreshold(alertThreshold);
|
||||
}
|
||||
|
||||
if (batchConflationEnabled != null) {
|
||||
gatewaySenderFactory.setBatchConflationEnabled(batchConflationEnabled);
|
||||
}
|
||||
|
||||
if (batchSize != null) {
|
||||
gatewaySenderFactory.setBatchSize(batchSize);
|
||||
}
|
||||
@@ -123,21 +134,15 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
|
||||
gatewaySenderFactory.setDispatcherThreads(dispatcherThreads);
|
||||
}
|
||||
|
||||
if (batchConflationEnabled != null) {
|
||||
gatewaySenderFactory.setBatchConflationEnabled(batchConflationEnabled);
|
||||
}
|
||||
|
||||
if (!CollectionUtils.isEmpty(eventFilters)) {
|
||||
for (GatewayEventFilter eventFilter : eventFilters) {
|
||||
gatewaySenderFactory.addGatewayEventFilter(eventFilter);
|
||||
}
|
||||
for (GatewayEventFilter eventFilter : CollectionUtils.nullSafeList(eventFilters)) {
|
||||
gatewaySenderFactory.addGatewayEventFilter(eventFilter);
|
||||
}
|
||||
|
||||
if (eventSubstitutionFilter != null) {
|
||||
gatewaySenderFactory.setGatewayEventSubstitutionFilter(eventSubstitutionFilter);
|
||||
}
|
||||
|
||||
gatewaySenderFactory.setManualStart(true);
|
||||
gatewaySenderFactory.setManualStart(manualStart);
|
||||
|
||||
if (maximumQueueMemory != null) {
|
||||
gatewaySenderFactory.setMaximumQueueMemory(maximumQueueMemory);
|
||||
@@ -163,10 +168,8 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
|
||||
gatewaySenderFactory.setSocketReadTimeout(socketReadTimeout);
|
||||
}
|
||||
|
||||
if (!CollectionUtils.isEmpty(transportFilters)) {
|
||||
for (GatewayTransportFilter transportFilter : transportFilters) {
|
||||
gatewaySenderFactory.addGatewayTransportFilter(transportFilter);
|
||||
}
|
||||
for (GatewayTransportFilter transportFilter : CollectionUtils.nullSafeList(transportFilters)) {
|
||||
gatewaySenderFactory.addGatewayTransportFilter(transportFilter);
|
||||
}
|
||||
|
||||
GatewaySenderWrapper wrapper = new GatewaySenderWrapper(gatewaySenderFactory.create(getName(),
|
||||
@@ -275,58 +278,4 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
|
||||
public void setTransportFilters(List<GatewayTransportFilter> gatewayTransportFilters) {
|
||||
this.transportFilters = gatewayTransportFilters;
|
||||
}
|
||||
|
||||
/* (non-Javadoc)
|
||||
* @see org.springframework.context.SmartLifecycle#isAutoStartup()
|
||||
*/
|
||||
@Override
|
||||
public boolean isAutoStartup() {
|
||||
return !manualStart;
|
||||
}
|
||||
|
||||
/* (non-Javadoc)
|
||||
* @see org.springframework.context.Phased#getPhase()
|
||||
*/
|
||||
@Override
|
||||
public int getPhase() {
|
||||
return Integer.MAX_VALUE;
|
||||
}
|
||||
|
||||
/* (non-Javadoc)
|
||||
* @see org.springframework.context.Lifecycle#isRunning()
|
||||
*/
|
||||
@Override
|
||||
public boolean isRunning() {
|
||||
return gatewaySender.isRunning();
|
||||
}
|
||||
|
||||
/* (non-Javadoc)
|
||||
* @see org.springframework.context.Lifecycle#start()
|
||||
*/
|
||||
@Override
|
||||
public synchronized void start() {
|
||||
Assert.notNull(gatewaySender, "The GatewaySender was not properly configured and initialized!");
|
||||
|
||||
if (!isRunning()){
|
||||
gatewaySender.start();
|
||||
}
|
||||
}
|
||||
|
||||
/* (non-Javadoc)
|
||||
* @see org.springframework.context.Lifecycle#stop()
|
||||
*/
|
||||
@Override
|
||||
public void stop() {
|
||||
gatewaySender.stop();
|
||||
}
|
||||
|
||||
/* (non-Javadoc)
|
||||
* @see org.springframework.context.SmartLifecycle#stop(java.lang.Runnable)
|
||||
*/
|
||||
@Override
|
||||
public void stop(Runnable callback) {
|
||||
stop();
|
||||
callback.run();
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -63,7 +63,6 @@ public class GatewayReceiverAutoStartNamespaceTest {
|
||||
public void testAuto() throws Exception {
|
||||
assertNotNull("The 'Auto' GatewayReceiverFactoryBean was not properly configured and initialized!",
|
||||
autoGatewayReceiverFactory);
|
||||
assertTrue(autoGatewayReceiverFactory.isAutoStartup());
|
||||
|
||||
GatewayReceiver autoGatewayReceiver = autoGatewayReceiverFactory.getObject();
|
||||
|
||||
|
||||
@@ -63,7 +63,6 @@ public class GatewayReceiverDefaultStartNamespaceTest {
|
||||
public void testDefault() throws Exception {
|
||||
assertNotNull("The 'Default' GatewayReceiverFactoryBean was not properly configured and initialized!",
|
||||
defaultGatewayReceiverFactory);
|
||||
assertTrue(defaultGatewayReceiverFactory.isAutoStartup());
|
||||
|
||||
GatewayReceiver defaultGatewayReceiver = defaultGatewayReceiverFactory.getObject();
|
||||
|
||||
|
||||
@@ -62,7 +62,6 @@ public class GatewayReceiverManualStartNamespaceTest {
|
||||
public void testManual() throws Exception {
|
||||
assertNotNull("The 'Manual' GatewayReceiverFactoryBean was not properly configured and initialized!",
|
||||
manualGatewayReceiverFactory);
|
||||
assertFalse(manualGatewayReceiverFactory.isAutoStartup());
|
||||
|
||||
GatewayReceiver manualGatewayReceiver = manualGatewayReceiverFactory.getObject();
|
||||
|
||||
@@ -75,5 +74,4 @@ public class GatewayReceiverManualStartNamespaceTest {
|
||||
assertFalse(manualGatewayReceiver.isRunning());
|
||||
assertEquals(8192, manualGatewayReceiver.getSocketBufferSize());
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -1,11 +1,11 @@
|
||||
/*
|
||||
* Copyright 2002-2013 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
|
||||
*
|
||||
*
|
||||
* http://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.
|
||||
@@ -142,6 +142,7 @@ public class StubGatewayReceiverFactory implements GatewayReceiverFactory {
|
||||
@Override
|
||||
public GatewayReceiverFactory setManualStart(final boolean manualStart) {
|
||||
this.manualStart = manualStart;
|
||||
this.running = !manualStart;
|
||||
return this;
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user