modified for Gemfire 7 WAN API changes

This commit is contained in:
David Turanski
2012-08-10 11:12:42 -04:00
parent 13836baee0
commit c2536e91d8
13 changed files with 319 additions and 47 deletions

View File

@@ -7,6 +7,7 @@
</configSuffixes>
<enableImports><![CDATA[false]]></enableImports>
<configs>
<config>src/test/resources/org/springframework/data/gemfire/support/GemfirePersistenceExceptionTranslationTest-context.xml</config>
</configs>
<configSets>
</configSets>

View File

@@ -8,7 +8,7 @@ slf4jVersion = 1.6.4
springVersion = 3.1.2.RELEASE
springDataCommonsVersion = 1.4.0.M1
#Temporary until release of 7.0
gemfireVersion = 7.0-beta
gemfireVersion = 7.0.Beta-SNAPSHOT
gemfireVersion6x = 6.6.3
# Testing

View File

@@ -1,6 +1,7 @@
Spring GemFire Samples
----------------------
NOTE: Spring Data GemFire Sample code has been moved to https://github.com/SpringSource/spring-gemfire-examples
-----------------------------------------------------------------------------------------------------------------
This folder contains various various demo applications and samples for Spring GemFire.
Please see each folder for detailed instructions (readme.txt).

View File

@@ -38,6 +38,7 @@ import org.springframework.util.Assert;
import org.springframework.util.ClassUtils;
import org.springframework.util.CollectionUtils;
import com.gemstone.gemfire.GemFireCheckedException;
import com.gemstone.gemfire.GemFireException;
import com.gemstone.gemfire.cache.Cache;
import com.gemstone.gemfire.cache.CacheClosedException;
@@ -454,6 +455,12 @@ public class CacheFactoryBean implements BeanNameAware, BeanFactoryAware, BeanCl
return wrapped;
}
}
if (ex.getCause() instanceof GemFireException) {
return GemfireCacheUtils.convertGemfireAccessException((GemFireException)ex.getCause());
}
if (ex.getCause() instanceof GemFireCheckedException) {
return GemfireCacheUtils.convertGemfireAccessException((GemFireCheckedException)ex.getCause());
}
return null;
}

View File

@@ -27,7 +27,7 @@ import org.springframework.util.Assert;
import org.springframework.util.ObjectUtils;
import org.springframework.util.ReflectionUtils;
import com.gemstone.gemfire.cache.AsyncEventQueue;
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;
@@ -133,13 +133,13 @@ public class RegionFactoryBean<K, V> extends RegionLookupFactoryBean<K, V> imple
hubId == null,
"It is invalid to configure a region with both a hubId and gatewaySenders. Note that the enableGateway and hubId properties are deprecated since Gemfire 7.0");
for (Object gatewaySender : gatewaySenders) {
regionFactory.addGatewaySender((GatewaySender) gatewaySender);
regionFactory.addGatewaySenderId(((GatewaySender) gatewaySender).getId());
}
}
if (!ObjectUtils.isEmpty(asyncEventQueues)) {
for (Object asyncEventQueue : asyncEventQueues) {
regionFactory.addAsyncEventQueue((AsyncEventQueue) asyncEventQueue);
regionFactory.addAsyncEventQueueId(((AsyncEventQueue) asyncEventQueue).getId());
}
}

View File

@@ -19,10 +19,8 @@ package org.springframework.data.gemfire.config;
import java.util.List;
import org.springframework.beans.factory.BeanDefinitionStoreException;
import org.springframework.beans.factory.config.BeanDefinition;
import org.springframework.beans.factory.support.AbstractBeanDefinition;
import org.springframework.beans.factory.support.BeanDefinitionBuilder;
import org.springframework.beans.factory.support.BeanDefinitionReaderUtils;
import org.springframework.beans.factory.support.ManagedList;
import org.springframework.beans.factory.support.ManagedMap;
import org.springframework.beans.factory.xml.AbstractSimpleBeanDefinitionParser;

View File

@@ -38,6 +38,8 @@ public abstract class AbstractWANComponentFactoryBean<T> implements FactoryBean<
protected final Cache cache;
protected Object factory;
protected AbstractWANComponentFactoryBean(Cache cache) {
this.cache = cache;
}
@@ -71,4 +73,9 @@ public abstract class AbstractWANComponentFactoryBean<T> implements FactoryBean<
public final boolean isSingleton() {
return true;
}
public void setFactory(Object factory) {
this.factory = factory;
}
}

View File

@@ -17,11 +17,11 @@ package org.springframework.data.gemfire.wan;
import org.springframework.util.Assert;
import com.gemstone.gemfire.cache.AsyncEventQueue;
import com.gemstone.gemfire.cache.AsyncEventQueueFactory;
import com.gemstone.gemfire.cache.asyncqueue.AsyncEventQueue;
import com.gemstone.gemfire.cache.asyncqueue.AsyncEventQueueFactory;
import com.gemstone.gemfire.cache.Cache;
import com.gemstone.gemfire.cache.CacheClosedException;
import com.gemstone.gemfire.cache.wan.AsyncEventListener;
import com.gemstone.gemfire.cache.asyncqueue.AsyncEventListener;
/**
* FactoryBean for creating GemFire {@link AsyncEventQueue}s.
@@ -44,8 +44,6 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
private Boolean persistent;
private Boolean parallel;
private String diskStoreRef;
/**
@@ -71,11 +69,13 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
@Override
protected void doInit() {
Assert.notNull(asyncEventListener, "AsyncEventListener cannot be null");
AsyncEventQueueFactory asyncEventQueueFactory = cache.createAsyncEventQueueFactory();
if (parallel != null) {
asyncEventQueueFactory.setParallel(parallel);
AsyncEventQueueFactory asyncEventQueueFactory = null;
if (this.factory == null) {
asyncEventQueueFactory = cache.createAsyncEventQueueFactory();
} else {
asyncEventQueueFactory = (AsyncEventQueueFactory)factory;
}
if (persistent != null) {
asyncEventQueueFactory.setPersistent(persistent);
}
@@ -99,7 +99,6 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
public void destroy() throws Exception {
if (!cache.isClosed()) {
try {
asyncEventQueue.stop();
asyncEventListener.close();
}
catch (CacheClosedException cce) {
@@ -119,8 +118,4 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
public void setPersistent(Boolean persistent) {
this.persistent = persistent;
}
public void setParallel(Boolean parallel) {
this.parallel = parallel;
}
}

View File

@@ -92,7 +92,12 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
@Override
protected void doInit() {
GatewaySenderFactory gatewaySenderFactory = cache.createGatewaySenderFactory();
GatewaySenderFactory gatewaySenderFactory = null;
if (this.factory == null) {
gatewaySenderFactory = cache.createGatewaySenderFactory();
} else {
gatewaySenderFactory = (GatewaySenderFactory)factory;
}
if (diskStoreRef != null) {
persistent = (persistent == null) ? Boolean.TRUE : persistent;
Assert.isTrue(persistent, "specifying a disk store requires persistent property to be true");
@@ -161,10 +166,6 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
gatewaySender = gatewaySenderFactory.create(name, remoteDistributedSystemId);
}
public void setGatewaySender(GatewaySender gatewaySender) {
this.gatewaySender = gatewaySender;
}
public void setRemoteDistributedSystemId(int remoteDistributedSystemId) {
this.remoteDistributedSystemId = remoteDistributedSystemId;
}

View File

@@ -29,14 +29,20 @@ import org.springframework.context.support.GenericXmlApplicationContext;
*/
public abstract class RecreatingContextTest {
protected GenericApplicationContext ctx;
protected GenericXmlApplicationContext ctx;
protected abstract String location();
protected void configureContext(){
}
@Before
public void createCtx() {
ctx = new GenericXmlApplicationContext(location());
ctx = new GenericXmlApplicationContext();
configureContext();
ctx.load(location());
ctx.registerShutdownHook();
ctx.refresh();
}
@After

View File

@@ -18,30 +18,39 @@ package org.springframework.data.gemfire.config;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertTrue;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
import java.io.File;
import java.io.FilenameFilter;
import java.io.InputStream;
import java.io.OutputStream;
import java.util.ArrayList;
import java.util.List;
import java.util.Set;
import org.junit.AfterClass;
import org.junit.Before;
import org.junit.Test;
import org.springframework.beans.BeansException;
import org.springframework.beans.factory.config.BeanPostProcessor;
import org.springframework.data.gemfire.RecreatingContextTest;
import org.springframework.data.gemfire.RegionFactoryBean;
import org.springframework.data.gemfire.TestUtils;
import org.springframework.data.gemfire.wan.AsyncEventQueueFactoryBean;
import org.springframework.data.gemfire.wan.GatewaySenderFactoryBean;
import com.gemstone.gemfire.cache.AsyncEventQueue;
import com.gemstone.gemfire.cache.Cache;
import com.gemstone.gemfire.cache.Region;
import com.gemstone.gemfire.cache.asyncqueue.AsyncEvent;
import com.gemstone.gemfire.cache.asyncqueue.AsyncEventListener;
import com.gemstone.gemfire.cache.asyncqueue.AsyncEventQueue;
import com.gemstone.gemfire.cache.asyncqueue.AsyncEventQueueFactory;
import com.gemstone.gemfire.cache.util.Gateway.OrderPolicy;
import com.gemstone.gemfire.cache.wan.AsyncEvent;
import com.gemstone.gemfire.cache.wan.AsyncEventListener;
import com.gemstone.gemfire.cache.wan.GatewayEventFilter;
import com.gemstone.gemfire.cache.wan.GatewayQueueEvent;
import com.gemstone.gemfire.cache.wan.GatewayReceiver;
import com.gemstone.gemfire.cache.wan.GatewaySender;
import com.gemstone.gemfire.cache.wan.GatewaySenderFactory;
import com.gemstone.gemfire.cache.wan.GatewayTransportFilter;
/**
@@ -50,11 +59,17 @@ import com.gemstone.gemfire.cache.wan.GatewayTransportFilter;
* @author David Turanski
*
*/
public class GemfireV7GatewayNamespaceTest extends RecreatingContextTest {
public class GemfireV7GatewayNamespaceTest extends RecreatingContextTest implements BeanPostProcessor {
@Override
protected String location() {
return "/org/springframework/data/gemfire/config/gateway-v7-ns.xml";
}
@Override
protected void configureContext() {
ctx.getBeanFactory().addBeanPostProcessor(this);
}
@Before
@Override
@@ -122,18 +137,21 @@ public class GemfireV7GatewayNamespaceTest extends RecreatingContextTest {
assertEquals(true, TestUtils.readField("manualStart", gwsfb));
}
private void testInnerGatewaySender() {
@SuppressWarnings("rawtypes")
private void testInnerGatewaySender() throws Exception {
Region<?, ?> region = ctx.getBean("region-inner-gateway-sender", Region.class);
GatewaySender gws = ctx.getBean("gateway-sender", GatewaySender.class);
assertNotNull(region.getAttributes().getGatewaySenders());
assertEquals(2, region.getAttributes().getGatewaySenders().size());
assertNotNull(region.getAttributes().getGatewaySenderIds());
assertEquals(2, region.getAttributes().getGatewaySenderIds().size());
// Isolate the inner gateway
Set<GatewaySender> gatewaySenders = region.getAttributes().getGatewaySenders();
assertTrue(gatewaySenders.remove(gws));
gatewaySenders.remove(gws);
// // Isolate the inner gateway
// Set<String> gatewaySenders = region.getAttributes().getGatewaySenderIds();
// assertTrue(gatewaySenders.remove(gws));
// gatewaySenders.remove(gws);
gws = gatewaySenders.iterator().next();
RegionFactoryBean rfb = ctx.getBean("&region-inner-gateway-sender", RegionFactoryBean.class);
Object[] gwsenders = TestUtils.readField("gatewaySenders", rfb);
gws = (GatewaySender)gwsenders[0];
List<GatewayEventFilter> eventFilters = gws.getGatewayEventFilters();
assertNotNull(eventFilters);
assertEquals(1, eventFilters.size());
@@ -179,18 +197,19 @@ public class GemfireV7GatewayNamespaceTest extends RecreatingContextTest {
}
@Override
public void afterAcknowledgement(AsyncEvent arg0) {
public void afterAcknowledgement(GatewayQueueEvent arg0) {
// TODO Auto-generated method stub
}
@Override
public boolean beforeEnqueue(AsyncEvent arg0) {
public boolean beforeEnqueue(GatewayQueueEvent arg0) {
// TODO Auto-generated method stub
return false;
}
@Override
public boolean beforeTransmit(AsyncEvent arg0) {
public boolean beforeTransmit(GatewayQueueEvent arg0) {
// TODO Auto-generated method stub
return false;
}
@@ -234,4 +253,242 @@ public class GemfireV7GatewayNamespaceTest extends RecreatingContextTest {
}
}
public static class StubAsyncEventQueueFactory implements AsyncEventQueueFactory {
private AsyncEventListener listener;
private AsyncEventQueue asyncEventQueue = mock(AsyncEventQueue.class);
private boolean persistent;
private int maxQueueMemory;
private String diskStoreName;
private int batchSize;
private String name;
@Override
public AsyncEventQueue create(String name, AsyncEventListener listener) {
this.name = name;
this.listener = listener;
when(asyncEventQueue.getAsyncEventListener()).thenReturn(this.listener);
when(asyncEventQueue.getBatchSize()).thenReturn(this.batchSize);
when(asyncEventQueue.getDiskStoreName()).thenReturn(this.diskStoreName);
when(asyncEventQueue.isPersistent()).thenReturn(this.persistent);
when(asyncEventQueue.getId()).thenReturn(this.name);
when(asyncEventQueue.getMaximumQueueMemory()).thenReturn(this.maxQueueMemory);
return this.asyncEventQueue;
}
@Override
public AsyncEventQueueFactory setBatchSize(int batchSize) {
this.batchSize = batchSize;
return this;
}
@Override
public AsyncEventQueueFactory setDiskStoreName(String diskStoreName) {
this.diskStoreName = diskStoreName;
return this;
}
@Override
public AsyncEventQueueFactory setMaximumQueueMemory(int maxQueueMemory) {
this.maxQueueMemory = maxQueueMemory;
return this;
}
@Override
public AsyncEventQueueFactory setPersistent(boolean persistent) {
this.persistent = persistent;
return this;
}
}
public static class StubGWSenderFactory implements GatewaySenderFactory {
private GatewaySender gatewaySender = mock(GatewaySender.class);
private int alertThreshold;
private boolean batchConflationEnabled;
private int batchSize;
private int batchTimeInterval;
private String diskStoreName;
private boolean diskSynchronous;
private int dispatcherThreads;
private boolean manualStart;
private int maxQueueMemory;
private OrderPolicy orderPolicy;
private boolean parallel;
private boolean persistenceEnabled;
private int socketBufferSize;
private int socketReadTimeout;
private List<GatewayEventFilter> eventFilters;
private List<GatewayTransportFilter> transportFilters;
private String name;
private int remoteSystemId;
public StubGWSenderFactory() {
this.eventFilters = new ArrayList<GatewayEventFilter>();
this.transportFilters = new ArrayList<GatewayTransportFilter>();
}
@Override
public GatewaySenderFactory addGatewayEventFilter(GatewayEventFilter filter) {
eventFilters.add(filter);
return this;
}
@Override
public GatewaySenderFactory addGatewayTransportFilter(GatewayTransportFilter filter) {
transportFilters.add(filter);
return this;
}
@Override
public GatewaySender create(String name, int remoteSystemId) {
this.name = name;
this.remoteSystemId = remoteSystemId;
when(gatewaySender.getId()).thenReturn(this.name);
when(gatewaySender.getRemoteDSId()).thenReturn(this.remoteSystemId);
when(gatewaySender.getAlertThreshold()).thenReturn(this.alertThreshold);
when(gatewaySender.getBatchSize()).thenReturn(this.batchSize);
when(gatewaySender.getBatchTimeInterval()).thenReturn(this.batchTimeInterval);
when(gatewaySender.getDiskStoreName()).thenReturn(this.diskStoreName);
when(gatewaySender.getDispatcherThreads()).thenReturn(this.dispatcherThreads);
when(gatewaySender.getGatewayEventFilters()).thenReturn(this.eventFilters);
when(gatewaySender.getGatewayTransportFilters()).thenReturn(this.transportFilters);
when(gatewaySender.getMaximumQueueMemory()).thenReturn(this.maxQueueMemory);
when(gatewaySender.getOrderPolicy()).thenReturn(this.orderPolicy);
when(gatewaySender.getSocketBufferSize()).thenReturn(this.socketBufferSize);
when(gatewaySender.getSocketReadTimeout()).thenReturn(this.socketReadTimeout);
when(gatewaySender.isManualStart()).thenReturn(this.manualStart);
when(gatewaySender.isBatchConflationEnabled()).thenReturn(this.batchConflationEnabled);
when(gatewaySender.isDiskSynchronous()).thenReturn(this.diskSynchronous);
when(gatewaySender.isParallel()).thenReturn(this.parallel);
when(gatewaySender.isPersistenceEnabled()).thenReturn(this.persistenceEnabled);
return gatewaySender;
}
@Override
public GatewaySenderFactory removeGatewayEventFilter(GatewayEventFilter filter) {
gatewaySender.removeGatewayEventFilter(filter);
return this;
}
@Override
public GatewaySenderFactory removeGatewayTransportFilter(GatewayTransportFilter filter) {
gatewaySender.removeGatewayTransportFilter(filter);
return this;
}
@Override
public GatewaySenderFactory setAlertThreshold(int alertThreshold) {
this.alertThreshold = alertThreshold;
return this;
}
@Override
public GatewaySenderFactory setBatchConflationEnabled(boolean batchConflationEnabled) {
this.batchConflationEnabled = batchConflationEnabled;
return this;
}
@Override
public GatewaySenderFactory setBatchSize(int batchSize) {
this.batchSize = batchSize;
return this;
}
@Override
public GatewaySenderFactory setBatchTimeInterval(int batchTimeInterval) {
this.batchTimeInterval = batchTimeInterval;
return this;
}
@Override
public GatewaySenderFactory setDiskStoreName(String diskStoreName) {
this.diskStoreName = diskStoreName;
return this;
}
@Override
public GatewaySenderFactory setDiskSynchronous(boolean diskSynchronous) {
this.diskSynchronous = diskSynchronous;
return this;
}
@Override
public GatewaySenderFactory setDispatcherThreads(int dispatcherThreads) {
this.dispatcherThreads = dispatcherThreads;
return this;
}
@Override
public GatewaySenderFactory setManualStart(boolean manualStart) {
this.manualStart = manualStart;
return this;
}
@Override
public GatewaySenderFactory setMaximumQueueMemory(int maxQueueMemory) {
this.maxQueueMemory = maxQueueMemory;
return this;
}
@Override
public GatewaySenderFactory setOrderPolicy(OrderPolicy orderPolicy) {
this.orderPolicy = orderPolicy;
return this;
}
@Override
public GatewaySenderFactory setParallel(boolean parallel) {
this.parallel = parallel;
return this;
}
@Override
public GatewaySenderFactory setPersistenceEnabled(boolean persistenceEnabled) {
this.persistenceEnabled = persistenceEnabled;
return this;
}
@Override
public GatewaySenderFactory setSocketBufferSize(int socketBufferSize) {
this.socketBufferSize = socketBufferSize;
return this;
}
@Override
public GatewaySenderFactory setSocketReadTimeout(int socketReadTimeout) {
this.socketReadTimeout = socketReadTimeout;
return this;
}
}
/*
* This mocks out the WAN components which are disabled in the developer edition
*/
@Override
public Object postProcessBeforeInitialization(Object bean, String beanName) throws BeansException {
if (bean instanceof GatewaySenderFactoryBean) {
((GatewaySenderFactoryBean)bean).setFactory(new StubGWSenderFactory());
}
if (bean instanceof AsyncEventQueueFactoryBean) {
((AsyncEventQueueFactoryBean)bean).setFactory(new StubAsyncEventQueueFactory());
}
return bean;
}
@Override
public Object postProcessAfterInitialization(Object bean, String beanName) throws BeansException {
return bean;
}
}

View File

@@ -6,7 +6,7 @@ log4j.appender.stdout.layout.ConversionPattern=%d %p [%c] - <%m>%n
log4j.category.org.springframework.data.gemfire.listener=TRACE
log4j.category.org.springframework.data.gemfire=DEBUG
log4j.category.org.springframework.beans=DEBUG
log4j.category.org.springframework=DEBUG
# for debugging datasource initialization
# log4j.category.test.jdbc=DEBUG

View File

@@ -10,11 +10,10 @@
<gfe:cache />
<!-- need manual-start=true for the unit test because GF will throw an exception if no locators are configured -->
<gfe:partitioned-region id="region-inner-gateway-sender" >
<gfe:gateway-sender
manual-start="true"
remote-distributed-system-id="1"
manual-start="true"
alert-threshold="10"
batch-size="11"
batch-time-interval="3000"