diff --git a/src/main/java/org/springframework/data/gemfire/RegionFactoryBean.java b/src/main/java/org/springframework/data/gemfire/RegionFactoryBean.java index ecceee72..878f2bc8 100644 --- a/src/main/java/org/springframework/data/gemfire/RegionFactoryBean.java +++ b/src/main/java/org/springframework/data/gemfire/RegionFactoryBean.java @@ -21,13 +21,14 @@ import java.lang.reflect.Field; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.springframework.beans.factory.DisposableBean; +import org.springframework.context.SmartLifecycle; import org.springframework.core.io.Resource; import org.springframework.data.gemfire.client.ClientRegionFactoryBean; +import org.springframework.data.gemfire.wan.SmartLifecycleGatewaySender; import org.springframework.util.Assert; import org.springframework.util.ObjectUtils; import org.springframework.util.ReflectionUtils; -import com.gemstone.gemfire.cache.asyncqueue.AsyncEventQueue; import com.gemstone.gemfire.cache.AttributesFactory; import com.gemstone.gemfire.cache.Cache; import com.gemstone.gemfire.cache.CacheClosedException; @@ -40,6 +41,7 @@ import com.gemstone.gemfire.cache.Region; import com.gemstone.gemfire.cache.RegionAttributes; import com.gemstone.gemfire.cache.RegionFactory; import com.gemstone.gemfire.cache.Scope; +import com.gemstone.gemfire.cache.asyncqueue.AsyncEventQueue; import com.gemstone.gemfire.cache.wan.GatewaySender; /** @@ -55,7 +57,11 @@ import com.gemstone.gemfire.cache.wan.GatewaySender; * @author Costin Leau * @author David Turanski */ -public class RegionFactoryBean extends RegionLookupFactoryBean implements DisposableBean { +public class RegionFactoryBean extends RegionLookupFactoryBean implements DisposableBean, SmartLifecycle { + + private boolean autoStartup = true; + + private boolean running; protected final Log log = LogFactory.getLog(getClass()); @@ -150,13 +156,13 @@ public class RegionFactoryBean extends RegionLookupFactoryBean imple if (cacheWriter != null) { regionFactory.setCacheWriter(cacheWriter); } - + if (diskStoreName != null) { regionFactory.setDiskStoreName(diskStoreName); - Assert.isTrue(!isNotPersistent(),"it is invalid to specify a disk store if 'persistent' is set to false."); + Assert.isTrue(!isNotPersistent(), "it is invalid to specify a disk store if 'persistent' is set to false."); persistent = true; } - + resolveDataPolicy(regionFactory, persistent, dataPolicy); if (scope != null) { @@ -194,8 +200,7 @@ public class RegionFactoryBean extends RegionLookupFactoryBean imple if (dataPolicy == null) { if (isPersistent()) { regionFactory.setDataPolicy(DataPolicy.PERSISTENT_REPLICATE); - } - else { + } else { regionFactory.setDataPolicy(DataPolicy.DEFAULT); } return; @@ -238,23 +243,21 @@ public class RegionFactoryBean extends RegionLookupFactoryBean imple * @param region */ protected void postProcess(Region region) { - // do nothing + } @Override public void destroy() throws Exception { if (region != null) { if (close) { - if (!region.getCache().isClosed()) { + if (!region.getRegionService().isClosed()) { try { region.close(); - } - catch (CacheClosedException cce) { + } catch (CacheClosedException cce) { // nothing to see folks, move on. } } - } - else if (destroy) { + } else if (destroy) { region.destroyRegion(); } } @@ -414,4 +417,73 @@ public class RegionFactoryBean extends RegionLookupFactoryBean imple protected boolean isNotPersistent() { return persistent != null && !persistent; } + + /* (non-Javadoc) + * @see org.springframework.context.Lifecycle#start() + */ + @Override + public void start() { + + if (!ObjectUtils.isEmpty(gatewaySenders)) { + synchronized (gatewaySenders) { + for (Object obj : gatewaySenders) { + SmartLifecycleGatewaySender gws = (SmartLifecycleGatewaySender) obj; + if (gws.isAutoStartup() && !gws.isRunning()) { + gws.start(); + } + } + } + } + this.running = true; + } + + /* (non-Javadoc) + * @see org.springframework.context.Lifecycle#stop() + */ + @Override + public void stop() { + if (!ObjectUtils.isEmpty(gatewaySenders)) { + synchronized (gatewaySenders) { + for (Object obj : gatewaySenders) { + SmartLifecycleGatewaySender gws = (SmartLifecycleGatewaySender) obj; + gws.stop(); + } + } + } + this.running = false; + } + + /* (non-Javadoc) + * @see org.springframework.context.Lifecycle#isRunning() + */ + @Override + public boolean isRunning() { + return this.running; + } + + /* (non-Javadoc) + * @see org.springframework.context.Phased#getPhase() + */ + @Override + public int getPhase() { + return Integer.MAX_VALUE; + } + + /* (non-Javadoc) + * @see org.springframework.context.SmartLifecycle#isAutoStartup() + */ + @Override + public boolean isAutoStartup() { + return this.autoStartup; + } + + /* (non-Javadoc) + * @see org.springframework.context.SmartLifecycle#stop(java.lang.Runnable) + */ + @Override + public void stop(Runnable callback) { + stop(); + callback.run(); + } + } diff --git a/src/main/java/org/springframework/data/gemfire/config/AsyncEventQueueParser.java b/src/main/java/org/springframework/data/gemfire/config/AsyncEventQueueParser.java index ef7d697c..7a29287c 100644 --- a/src/main/java/org/springframework/data/gemfire/config/AsyncEventQueueParser.java +++ b/src/main/java/org/springframework/data/gemfire/config/AsyncEventQueueParser.java @@ -19,11 +19,10 @@ import org.springframework.beans.factory.support.BeanDefinitionBuilder; import org.springframework.beans.factory.xml.AbstractSingleBeanDefinitionParser; import org.springframework.beans.factory.xml.ParserContext; import org.springframework.data.gemfire.wan.AsyncEventQueueFactoryBean; +import org.springframework.util.StringUtils; import org.springframework.util.xml.DomUtils; import org.w3c.dom.Element; -import com.gemstone.gemfire.internal.lang.StringUtils; - /** * @author David Turanski * @@ -41,8 +40,10 @@ public class AsyncEventQueueParser extends AbstractSingleBeanDefinitionParser { Element asyncEventListenerElement = DomUtils.getChildElementByTagName(element, "async-event-listener"); Object asyncEventListener = ParsingUtils.parseRefOrSingleNestedBeanDeclaration(parserContext, asyncEventListenerElement, builder); - String cacheName = StringUtils.isEmpty(element.getAttribute("cache-ref")) ? "gemfireCache" : element - .getAttribute("cache-ref"); + + String cacheName = !StringUtils.hasText(element.getAttribute("cache-ref")) ? GemfireConstants.DEFAULT_GEMFIRE_CACHE_NAME + : element.getAttribute("cache-ref"); + builder.addConstructorArgReference(cacheName); builder.addConstructorArgValue(asyncEventListener); ParsingUtils.setPropertyValue(element, builder, "batch-size"); @@ -50,5 +51,23 @@ public class AsyncEventQueueParser extends AbstractSingleBeanDefinitionParser { ParsingUtils.setPropertyValue(element, builder, "disk-store-ref"); ParsingUtils.setPropertyValue(element, builder, "persistent"); ParsingUtils.setPropertyValue(element, builder, "parallel"); + ParsingUtils.setPropertyValue(element, builder, NAME_ATTRIBUTE); + + if (!StringUtils.hasText(element.getAttribute(NAME_ATTRIBUTE))) { + if (element.getParentNode().getNodeName().endsWith("region")) { + Element region = (Element) element.getParentNode(); + String regionName = StringUtils.hasText(region.getAttribute("name")) ? region.getAttribute("name") + : region.getAttribute("id"); + + int i = 0; + String name = regionName + ".asyncEventQueue#" + i; + while (parserContext.getRegistry().isBeanNameInUse(name)) { + i++; + name = regionName + ".asyncEventQueue#" + i; + } + + builder.addPropertyValue("name", name); + } + } } } diff --git a/src/main/java/org/springframework/data/gemfire/config/GatewaySenderParser.java b/src/main/java/org/springframework/data/gemfire/config/GatewaySenderParser.java index a42392b0..e79880ac 100644 --- a/src/main/java/org/springframework/data/gemfire/config/GatewaySenderParser.java +++ b/src/main/java/org/springframework/data/gemfire/config/GatewaySenderParser.java @@ -36,10 +36,12 @@ class GatewaySenderParser extends AbstractSimpleBeanDefinitionParser { @Override protected void doParse(Element element, ParserContext parserContext, BeanDefinitionBuilder builder) { - builder.setLazyInit(false); + String cacheRef = element.getAttribute("cache-ref"); + // add cache reference (fallback to default if nothing is specified) - builder.addConstructorArgReference((StringUtils.hasText(cacheRef) ? cacheRef : "gemfireCache")); + builder.addConstructorArgReference((StringUtils.hasText(cacheRef) ? cacheRef + : GemfireConstants.DEFAULT_GEMFIRE_CACHE_NAME)); ParsingUtils.setPropertyValue(element, builder, "alert-threshold"); ParsingUtils.setPropertyValue(element, builder, "batch-size"); ParsingUtils.setPropertyValue(element, builder, "batch-time-interval"); @@ -55,6 +57,7 @@ class GatewaySenderParser extends AbstractSimpleBeanDefinitionParser { ParsingUtils.setPropertyValue(element, builder, "socket-read-timeout"); ParsingUtils.setPropertyValue(element, builder, "persistent"); ParsingUtils.setPropertyValue(element, builder, "parallel"); + ParsingUtils.setPropertyValue(element, builder, NAME_ATTRIBUTE); Element eventFilterElement = DomUtils.getChildElementByTagName(element, "event-filter"); if (eventFilterElement != null) { @@ -63,5 +66,25 @@ class GatewaySenderParser extends AbstractSimpleBeanDefinitionParser { } ParsingUtils.parseTransportFilters(element, parserContext, builder); + + /** + * set the name for an inner bean + */ + if (!StringUtils.hasText(element.getAttribute(NAME_ATTRIBUTE))) { + if (element.getParentNode().getNodeName().endsWith("region")) { + Element region = (Element) element.getParentNode(); + String regionName = StringUtils.hasText(region.getAttribute("name")) ? region.getAttribute("name") + : region.getAttribute("id"); + + int i = 0; + String name = regionName + ".gatewaySender#" + i; + while (parserContext.getRegistry().isBeanNameInUse(name)) { + i++; + name = regionName + ".gatewaySender#" + i; + } + + builder.addPropertyValue("name", name); + } + } } } 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 9c3c3802..fbea4c96 100644 --- a/src/main/java/org/springframework/data/gemfire/wan/AbstractWANComponentFactoryBean.java +++ b/src/main/java/org/springframework/data/gemfire/wan/AbstractWANComponentFactoryBean.java @@ -34,15 +34,25 @@ public abstract class AbstractWANComponentFactoryBean implements FactoryBean< DisposableBean { protected Log log = LogFactory.getLog(this.getClass()); - protected String name; + private String name; protected final Cache cache; protected Object factory; + private String beanName; + protected AbstractWANComponentFactoryBean(Cache cache) { this.cache = cache; } + + public void setName(String name) { + this.name = name; + } + + public String getName() { + return name!=null ? name: beanName; + } @Override public void destroy() throws Exception { @@ -50,13 +60,13 @@ public abstract class AbstractWANComponentFactoryBean implements FactoryBean< } @Override - public final void setBeanName(String name) { - this.name = name; + public final void setBeanName(String beanName) { + this.beanName = beanName; } @Override public final void afterPropertiesSet() throws Exception { - Assert.notNull(name, "Name cannot be null"); + Assert.notNull(getName(), "Name cannot be null"); Assert.notNull(cache, "Cache cannot be null"); doInit(); } diff --git a/src/main/java/org/springframework/data/gemfire/wan/AsyncEventQueueFactoryBean.java b/src/main/java/org/springframework/data/gemfire/wan/AsyncEventQueueFactoryBean.java index 255b2381..1a6196c2 100644 --- a/src/main/java/org/springframework/data/gemfire/wan/AsyncEventQueueFactoryBean.java +++ b/src/main/java/org/springframework/data/gemfire/wan/AsyncEventQueueFactoryBean.java @@ -92,7 +92,7 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean< asyncEventQueueFactory.setMaximumQueueMemory(maximumQueueMemory); } - asyncEventQueue = asyncEventQueueFactory.create(name, asyncEventListener); + asyncEventQueue = asyncEventQueueFactory.create(getName(), asyncEventListener); } @Override diff --git a/src/main/java/org/springframework/data/gemfire/wan/GatewayHubFactoryBean.java b/src/main/java/org/springframework/data/gemfire/wan/GatewayHubFactoryBean.java index ad914526..91515ff3 100644 --- a/src/main/java/org/springframework/data/gemfire/wan/GatewayHubFactoryBean.java +++ b/src/main/java/org/springframework/data/gemfire/wan/GatewayHubFactoryBean.java @@ -76,6 +76,7 @@ public class GatewayHubFactoryBean extends AbstractWANComponentFactoryBean { +public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean { private static List validOrderPolicyValues = Arrays.asList("KEY", "PARTITION", "THREAD"); private GatewaySender gatewaySender; @@ -58,7 +58,7 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean getObjectType() { - return GatewaySender.class; + return SmartLifecycleGatewaySender.class; } @Override @@ -151,9 +151,9 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean getGatewayEventFilters() { + return this.delegate.getGatewayEventFilters(); + } + + /* (non-Javadoc) + * @see com.gemstone.gemfire.cache.wan.GatewaySender#getGatewayTransportFilters() + */ + @Override + public List getGatewayTransportFilters() { + return this.delegate.getGatewayTransportFilters(); + } + + /* (non-Javadoc) + * @see com.gemstone.gemfire.cache.wan.GatewaySender#getId() + */ + @Override + public String getId() { + return this.delegate.getId(); + } + + /* (non-Javadoc) + * @see com.gemstone.gemfire.cache.wan.GatewaySender#getMaximumQueueMemory() + */ + @Override + public int getMaximumQueueMemory() { + return this.delegate.getMaximumQueueMemory(); + } + + /* (non-Javadoc) + * @see com.gemstone.gemfire.cache.wan.GatewaySender#getOrderPolicy() + */ + @Override + public OrderPolicy getOrderPolicy() { + return this.delegate.getOrderPolicy(); + } + + /* (non-Javadoc) + * @see com.gemstone.gemfire.cache.wan.GatewaySender#getRemoteDSId() + */ + @Override + public int getRemoteDSId() { + return this.delegate.getRemoteDSId(); + } + + /* (non-Javadoc) + * @see com.gemstone.gemfire.cache.wan.GatewaySender#getSocketBufferSize() + */ + @Override + public int getSocketBufferSize() { + return this.delegate.getSocketBufferSize(); + } + + /* (non-Javadoc) + * @see com.gemstone.gemfire.cache.wan.GatewaySender#getSocketReadTimeout() + */ + @Override + public int getSocketReadTimeout() { + return this.delegate.getSocketReadTimeout(); + } + + /* (non-Javadoc) + * @see com.gemstone.gemfire.cache.wan.GatewaySender#isBatchConflationEnabled() + */ + @Override + public boolean isBatchConflationEnabled() { + return this.delegate.isBatchConflationEnabled(); + } + + /* (non-Javadoc) + * @see com.gemstone.gemfire.cache.wan.GatewaySender#isDiskSynchronous() + */ + @Override + public boolean isDiskSynchronous() { + return this.delegate.isDiskSynchronous(); + } + + /* (non-Javadoc) + * @see com.gemstone.gemfire.cache.wan.GatewaySender#isManualStart() + */ + @Override + public boolean isManualStart() { + return true; + } + + /* (non-Javadoc) + * @see com.gemstone.gemfire.cache.wan.GatewaySender#isParallel() + */ + @Override + public boolean isParallel() { + return this.delegate.isParallel(); + } + + /* (non-Javadoc) + * @see com.gemstone.gemfire.cache.wan.GatewaySender#isPaused() + */ + @Override + public boolean isPaused() { + return this.delegate.isPaused(); + } + + /* (non-Javadoc) + * @see com.gemstone.gemfire.cache.wan.GatewaySender#isPersistenceEnabled() + */ + @Override + public boolean isPersistenceEnabled() { + return this.delegate.isPersistenceEnabled(); + } + + /* (non-Javadoc) + * @see com.gemstone.gemfire.cache.wan.GatewaySender#isRunning() + */ + @Override + public boolean isRunning() { + return this.delegate.isRunning(); + } + + /* (non-Javadoc) + * @see com.gemstone.gemfire.cache.wan.GatewaySender#pause() + */ + @Override + public void pause() { + this.delegate.pause(); + } + + /* (non-Javadoc) + * @see com.gemstone.gemfire.cache.wan.GatewaySender#removeGatewayEventFilter(com.gemstone.gemfire.cache.wan.GatewayEventFilter) + */ + @Override + public void removeGatewayEventFilter(GatewayEventFilter eventFilter) { + this.delegate.removeGatewayEventFilter(eventFilter); + } + + /* (non-Javadoc) + * @see com.gemstone.gemfire.cache.wan.GatewaySender#resume() + */ + @Override + public void resume() { + this.delegate.resume(); + + } + + /* (non-Javadoc) + * @see com.gemstone.gemfire.cache.wan.GatewaySender#start() + */ + @Override + public void start() { + this.delegate.start(); + + } + + /* (non-Javadoc) + * @see com.gemstone.gemfire.cache.wan.GatewaySender#stop() + */ + @Override + public void stop() { + this.delegate.stop(); + } + + /* (non-Javadoc) + * @see org.springframework.context.Phased#getPhase() + */ + @Override + public int getPhase() { + return Integer.MAX_VALUE; + } + + /* (non-Javadoc) + * @see org.springframework.context.SmartLifecycle#isAutoStartup() + */ + @Override + public boolean isAutoStartup() { + return this.autoStartup ; + } + + /* (non-Javadoc) + * @see org.springframework.context.SmartLifecycle#stop(java.lang.Runnable) + */ + @Override + public void stop(Runnable callback) { + stop(); + callback.run(); + } + +} diff --git a/src/main/resources/org/springframework/data/gemfire/config/spring-gemfire-1.2.xsd b/src/main/resources/org/springframework/data/gemfire/config/spring-gemfire-1.2.xsd index 455c9205..188f117d 100755 --- a/src/main/resources/org/springframework/data/gemfire/config/spring-gemfire-1.2.xsd +++ b/src/main/resources/org/springframework/data/gemfire/config/spring-gemfire-1.2.xsd @@ -2258,6 +2258,13 @@ Inner bean definition of the event filter minOccurs="0" maxOccurs="1" /> + + + + + @@ -2506,6 +2513,13 @@ use inner bean declarations. + + + + +