Fixing test failures, additional error handling
This commit is contained in:
@@ -168,6 +168,13 @@
|
||||
<version>${org.springframework.version}</version>
|
||||
</dependency>
|
||||
|
||||
<!-- Groovy -->
|
||||
<dependency>
|
||||
<groupId>org.codehaus.groovy</groupId>
|
||||
<artifactId>groovy-all</artifactId>
|
||||
<version>1.7.5</version>
|
||||
</dependency>
|
||||
|
||||
<!-- Spring Data -->
|
||||
<dependency>
|
||||
<groupId>org.springframework.data</groupId>
|
||||
|
||||
@@ -97,6 +97,12 @@
|
||||
<scope>test</scope>
|
||||
</dependency>
|
||||
|
||||
<!-- Groovy -->
|
||||
<dependency>
|
||||
<groupId>org.codehaus.groovy</groupId>
|
||||
<artifactId>groovy-all</artifactId>
|
||||
</dependency>
|
||||
|
||||
<dependency>
|
||||
<groupId>junit</groupId>
|
||||
<artifactId>junit</artifactId>
|
||||
|
||||
@@ -28,6 +28,13 @@ public abstract class AbstractAsyncOperation<T> implements Callable<T>, Initiali
|
||||
|
||||
protected RiakTemplate riakTemplate;
|
||||
|
||||
protected AbstractAsyncOperation() {
|
||||
}
|
||||
|
||||
protected AbstractAsyncOperation(RiakTemplate riakTemplate) {
|
||||
this.riakTemplate = riakTemplate;
|
||||
}
|
||||
|
||||
public RiakTemplate getRiakTemplate() {
|
||||
return riakTemplate;
|
||||
}
|
||||
|
||||
@@ -16,17 +16,27 @@
|
||||
|
||||
package org.springframework.datastore.riak.core;
|
||||
|
||||
import org.codehaus.groovy.runtime.GStringImpl;
|
||||
import org.codehaus.jackson.map.ObjectMapper;
|
||||
import org.codehaus.jackson.map.ser.CustomSerializerFactory;
|
||||
import org.codehaus.jackson.map.ser.ToStringSerializer;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.beans.factory.InitializingBean;
|
||||
import org.springframework.core.convert.ConversionService;
|
||||
import org.springframework.core.convert.support.ConversionServiceFactory;
|
||||
import org.springframework.dao.DataAccessResourceFailureException;
|
||||
import org.springframework.datastore.riak.convert.KeyValueStoreMetaData;
|
||||
import org.springframework.http.HttpEntity;
|
||||
import org.springframework.http.HttpHeaders;
|
||||
import org.springframework.http.HttpStatus;
|
||||
import org.springframework.http.MediaType;
|
||||
import org.springframework.http.client.ClientHttpRequestFactory;
|
||||
import org.springframework.http.converter.HttpMessageConverter;
|
||||
import org.springframework.http.converter.json.MappingJacksonHttpMessageConverter;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.ClassUtils;
|
||||
import org.springframework.web.client.HttpClientErrorException;
|
||||
import org.springframework.web.client.ResourceAccessException;
|
||||
import org.springframework.web.client.RestTemplate;
|
||||
import org.springframework.web.client.support.RestGatewaySupport;
|
||||
@@ -34,6 +44,7 @@ import org.springframework.web.client.support.RestGatewaySupport;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ConcurrentSkipListMap;
|
||||
|
||||
/**
|
||||
* @author J. Brisbin <jon@jbrisbin.com>
|
||||
@@ -41,8 +52,11 @@ import java.util.Map;
|
||||
@SuppressWarnings({"unchecked"})
|
||||
public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOperations, InitializingBean {
|
||||
|
||||
private static final boolean groovyPresent = ClassUtils.isPresent("org.codehaus.groovy.runtime.GStringImpl",
|
||||
RiakTemplate.class.getClassLoader());
|
||||
protected final Logger log = LoggerFactory.getLogger(getClass());
|
||||
protected ConversionService conversionService = ConversionServiceFactory.createDefaultConversionService();
|
||||
protected ConcurrentSkipListMap<Object, Object> cache = new ConcurrentSkipListMap<Object, Object>();
|
||||
protected String defaultUri = "http://localhost:8098/riak/{bucket}/{key}";
|
||||
|
||||
public RiakTemplate() {
|
||||
@@ -71,10 +85,7 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
|
||||
}
|
||||
|
||||
public <V> KeyValueStoreOperations set(Object key, V value) {
|
||||
String[] bucketAndKey = getBucketAndKey(key);
|
||||
if (null == bucketAndKey[0]) {
|
||||
bucketAndKey[0] = value.getClass().getName();
|
||||
}
|
||||
String[] bucketAndKey = extractBucketAndKey(key);
|
||||
if (null == bucketAndKey[1]) {
|
||||
// TODO: Handle auto-generation of key name
|
||||
}
|
||||
@@ -91,7 +102,7 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
|
||||
}
|
||||
|
||||
public KeyValueStoreOperations setAsBytes(Object key, byte[] value) {
|
||||
String[] bucketAndKey = getBucketAndKey(key);
|
||||
String[] bucketAndKey = extractBucketAndKey(key);
|
||||
if (null == bucketAndKey[0]) {
|
||||
bucketAndKey[0] = "bytes";
|
||||
}
|
||||
@@ -111,7 +122,7 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
|
||||
}
|
||||
|
||||
public <V> V get(Object key) {
|
||||
String[] bucketAndKey = getBucketAndKey(key);
|
||||
String[] bucketAndKey = extractBucketAndKey(key);
|
||||
Assert.noNullElements(bucketAndKey, "Must specify a bucket and key to retrieve.");
|
||||
RestTemplate restTemplate = getRestTemplate();
|
||||
Class targetClass;
|
||||
@@ -123,7 +134,14 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
|
||||
if (log.isDebugEnabled()) {
|
||||
log.debug(String.format("GET object: key=%s", key));
|
||||
}
|
||||
return (V) restTemplate.getForObject(defaultUri, targetClass, (Object[]) bucketAndKey);
|
||||
try {
|
||||
return (V) restTemplate.getForObject(defaultUri, targetClass, (Object[]) bucketAndKey);
|
||||
} catch (HttpClientErrorException e) {
|
||||
if (e.getStatusCode() != HttpStatus.NOT_FOUND) {
|
||||
throw new DataAccessResourceFailureException(e.getMessage(), e);
|
||||
}
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
public byte[] getAsBytes(Object key) {
|
||||
@@ -131,7 +149,7 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
|
||||
}
|
||||
|
||||
public <T> T getAsType(Object key, Class<T> requiredType) {
|
||||
String[] bucketAndKey = getBucketAndKey(key);
|
||||
String[] bucketAndKey = extractBucketAndKey(key);
|
||||
if (null == bucketAndKey[0]) {
|
||||
bucketAndKey[0] = requiredType.getName();
|
||||
}
|
||||
@@ -140,7 +158,14 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
|
||||
if (log.isDebugEnabled()) {
|
||||
log.debug(String.format("GET object: key=%s, type=%s", key, requiredType.getName()));
|
||||
}
|
||||
return (T) restTemplate.getForObject(defaultUri, requiredType, (Object[]) bucketAndKey);
|
||||
try {
|
||||
return (T) restTemplate.getForObject(defaultUri, requiredType, (Object[]) bucketAndKey);
|
||||
} catch (HttpClientErrorException e) {
|
||||
if (e.getStatusCode() != HttpStatus.NOT_FOUND) {
|
||||
throw new DataAccessResourceFailureException(e.getMessage(), e);
|
||||
}
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
public <V> V getAndSet(Object key, V value) {
|
||||
@@ -237,7 +262,7 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
|
||||
}
|
||||
|
||||
public boolean containsKey(Object key) {
|
||||
String[] bucketAndKey = getBucketAndKey(key);
|
||||
String[] bucketAndKey = extractBucketAndKey(key);
|
||||
Assert.noNullElements(bucketAndKey, "Must specify a bucket and key to check for.");
|
||||
RestTemplate restTemplate = getRestTemplate();
|
||||
HttpHeaders headers = null;
|
||||
@@ -249,22 +274,44 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
|
||||
}
|
||||
|
||||
public boolean deleteKeys(Object... keys) {
|
||||
boolean deleted = false;
|
||||
boolean stillExists = false;
|
||||
for (Object key : keys) {
|
||||
String[] bucketAndKey = getBucketAndKey(key);
|
||||
String[] bucketAndKey = extractBucketAndKey(key);
|
||||
Assert.noNullElements(bucketAndKey, "Must specify a bucket and key to delete.");
|
||||
RestTemplate restTemplate = getRestTemplate();
|
||||
restTemplate.delete(defaultUri, (Object[]) bucketAndKey);
|
||||
deleted = (!deleted && containsKey(key) ? false : true);
|
||||
try {
|
||||
restTemplate.delete(defaultUri, (Object[]) bucketAndKey);
|
||||
} catch (HttpClientErrorException e) {
|
||||
if (e.getStatusCode() != HttpStatus.NOT_FOUND) {
|
||||
throw new DataAccessResourceFailureException(e.getMessage(), e);
|
||||
}
|
||||
}
|
||||
if (!stillExists) {
|
||||
stillExists = containsKey(key);
|
||||
}
|
||||
}
|
||||
return deleted;
|
||||
return !stillExists;
|
||||
}
|
||||
|
||||
public void afterPropertiesSet() throws Exception {
|
||||
Assert.notNull(conversionService, "Must specify a valid ConversionService.");
|
||||
|
||||
if (groovyPresent) {
|
||||
// Native conversion for Groovy GString objects
|
||||
List<HttpMessageConverter<?>> converters = getRestTemplate().getMessageConverters();
|
||||
for (HttpMessageConverter converter : converters) {
|
||||
if (converter instanceof MappingJacksonHttpMessageConverter) {
|
||||
ObjectMapper mapper = new ObjectMapper();
|
||||
CustomSerializerFactory fac = new CustomSerializerFactory();
|
||||
fac.addSpecificMapping(GStringImpl.class, ToStringSerializer.instance);
|
||||
mapper.setSerializerFactory(fac);
|
||||
((MappingJacksonHttpMessageConverter) converter).setObjectMapper(mapper);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
protected String[] getBucketAndKey(Object obj) {
|
||||
protected String[] extractBucketAndKey(Object obj) {
|
||||
Object bucket = null;
|
||||
Object key = null;
|
||||
if (obj instanceof Map) {
|
||||
@@ -272,6 +319,11 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
|
||||
bucket = m.get("bucket");
|
||||
key = m.get("key");
|
||||
} else {
|
||||
// Override from Annotation?
|
||||
KeyValueStoreMetaData meta = obj.getClass().getAnnotation(KeyValueStoreMetaData.class);
|
||||
if (null != meta && null != meta.family()) {
|
||||
bucket = meta.family();
|
||||
}
|
||||
String s = obj.toString();
|
||||
if (s.contains("@")) {
|
||||
// This is likely the result of Object.toString()
|
||||
@@ -281,12 +333,17 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
|
||||
}
|
||||
if (s.contains(":")) {
|
||||
String[] a = s.split(":");
|
||||
bucket = a[0];
|
||||
if (null == bucket) {
|
||||
bucket = a[0];
|
||||
}
|
||||
key = a[1];
|
||||
} else {
|
||||
bucket = null;
|
||||
key = s;
|
||||
}
|
||||
if (null == bucket) {
|
||||
// Default to the class name for the bucket
|
||||
bucket = (obj.getClass() == byte[].class ? "bytes" : obj.getClass().getName());
|
||||
}
|
||||
}
|
||||
return new String[]{(null != bucket ? bucket.toString() : null), (null != key ? key.toString() : null)};
|
||||
}
|
||||
|
||||
@@ -14,15 +14,15 @@
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.springframework.datastore.riak.convert;
|
||||
|
||||
import org.springframework.core.convert.support.GenericConversionService;
|
||||
package org.springframework.datastore.riak.mapreduce;
|
||||
|
||||
/**
|
||||
* @author J. Brisbin <jon@jbrisbin.com>
|
||||
*/
|
||||
public class RiakConversionService extends GenericConversionService{
|
||||
public interface MapReduceOperation {
|
||||
|
||||
String getType();
|
||||
|
||||
Object getRepresentation();
|
||||
|
||||
public RiakConversionService() {
|
||||
}
|
||||
}
|
||||
@@ -36,14 +36,15 @@ class RiakTemplateSpec extends Specification {
|
||||
|
||||
given:
|
||||
def i = run++
|
||||
def objIn = [test: "value $i".toString(), integer: 12]
|
||||
String val = "value $i"
|
||||
def objIn = [test: "value $i", integer: 12]
|
||||
riak.set("test:test", objIn)
|
||||
|
||||
when:
|
||||
def objOut = riak.get("test:test")
|
||||
|
||||
then:
|
||||
objOut.test == "value $i"
|
||||
objOut.test == val
|
||||
|
||||
}
|
||||
|
||||
@@ -51,14 +52,15 @@ class RiakTemplateSpec extends Specification {
|
||||
|
||||
given:
|
||||
def i = run++
|
||||
def objIn = [test: "value $i".toString(), integer: 12]
|
||||
String val = "value $i"
|
||||
def objIn = [test: val, integer: 12]
|
||||
riak.set([bucket: "test", key: "test"], objIn)
|
||||
|
||||
when:
|
||||
def objOut = riak.get([bucket: "test", key: "test"])
|
||||
|
||||
then:
|
||||
objOut.test == "value $i"
|
||||
objOut.test == val
|
||||
|
||||
}
|
||||
|
||||
@@ -103,7 +105,7 @@ class RiakTemplateSpec extends Specification {
|
||||
def "Test multiple get"() {
|
||||
|
||||
when:
|
||||
def objs = riak.getValues(["test:test", "${TestObject.name}:test".toString()])
|
||||
def objs = riak.getValues(["test:test", "${TestObject.name}:test"])
|
||||
|
||||
then:
|
||||
2 == objs.size()
|
||||
@@ -114,7 +116,8 @@ class RiakTemplateSpec extends Specification {
|
||||
|
||||
given:
|
||||
def i = run++
|
||||
def newObj = [test: "value $i".toString(), integer: 12]
|
||||
String val = "value $i"
|
||||
def newObj = [test: val, integer: 12]
|
||||
|
||||
when:
|
||||
def oldObj = riak.getAndSet("test:test", newObj)
|
||||
@@ -124,33 +127,35 @@ class RiakTemplateSpec extends Specification {
|
||||
|
||||
}
|
||||
|
||||
def "Test setMultipleIfKeysNonExistent with Map"() {
|
||||
|
||||
given:
|
||||
def i = run++
|
||||
String firstKey = "test:test$i"
|
||||
String secondKey = "${TestObject.name}:test$i"
|
||||
def newObj = [
|
||||
"$firstKey": [test: "value $i".toString(), integer: 12],
|
||||
"$secondKey": [test: "value $i".toString(), integer: 12]
|
||||
]
|
||||
|
||||
when:
|
||||
def secondObj = riak.setMultipleIfKeysNonExistent(newObj).get(secondKey)
|
||||
|
||||
then:
|
||||
"value $i" == secondObj.test
|
||||
|
||||
}
|
||||
|
||||
def "Test deleteKeys"() {
|
||||
|
||||
when:
|
||||
def deleted = riak.deleteKeys("test:test", "${TestObject.name}:test".toString())
|
||||
def deleted = riak.deleteKeys("test:test", "${TestObject.name}:test")
|
||||
|
||||
then:
|
||||
true == deleted
|
||||
|
||||
}
|
||||
|
||||
def "Test setMultipleIfKeysNonExistent with Map"() {
|
||||
|
||||
given:
|
||||
def newObj = [
|
||||
"test:test": [test: "value", integer: 12],
|
||||
"${TestObject.name}:test": [test: "value", integer: 12]
|
||||
]
|
||||
|
||||
when:
|
||||
def secondObj = riak.setMultipleIfKeysNonExistent(newObj).get("${TestObject.name}:test")
|
||||
secondObj.test = "newValue"
|
||||
def thirdObj = riak.setMultipleIfKeysNonExistent(["${TestObject.name}:test": secondObj]).get("${TestObject.name}:test")
|
||||
|
||||
then:
|
||||
"value" == thirdObj.test
|
||||
|
||||
cleanup:
|
||||
riak.deleteKeys("test:test", "${TestObject.name}:test")
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user