Fix new Sonar smells for MMIH (#2795)
* Fix new Sonar smells for MMIH * Fix complexity smells in the `MessagingMethodInvokerHelper` * Reuse `get()` in the `ZookeeperMetadataStore` in case of `putIfAbsent()` error * Polishing for `ZookeeperMetadataStoreTests` * Fix for `hasUnqualifiedMapParameter` * Fix `findHandlerMethodsForTarget()` complexity * * Fix typo and remove FQCN for `ClassUtils`
This commit is contained in:
committed by
Gary Russell
parent
7f806bd94c
commit
b4273dd662
@@ -22,7 +22,6 @@ import java.lang.reflect.Modifier;
|
||||
import java.lang.reflect.ParameterizedType;
|
||||
import java.lang.reflect.Proxy;
|
||||
import java.lang.reflect.Type;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.Collection;
|
||||
import java.util.Collections;
|
||||
@@ -80,7 +79,6 @@ import org.springframework.integration.support.json.JsonObjectMapper;
|
||||
import org.springframework.integration.support.json.JsonObjectMapperProvider;
|
||||
import org.springframework.integration.util.AbstractExpressionEvaluator;
|
||||
import org.springframework.integration.util.AnnotatedMethodFilter;
|
||||
import org.springframework.integration.util.ClassUtils;
|
||||
import org.springframework.integration.util.FixedMethodFilter;
|
||||
import org.springframework.integration.util.MessagingAnnotationUtils;
|
||||
import org.springframework.integration.util.UniqueMethodFilter;
|
||||
@@ -99,10 +97,10 @@ import org.springframework.messaging.handler.invocation.HandlerMethodArgumentRes
|
||||
import org.springframework.messaging.handler.invocation.InvocableHandlerMethod;
|
||||
import org.springframework.messaging.handler.invocation.MethodArgumentResolutionException;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.ClassUtils;
|
||||
import org.springframework.util.CollectionUtils;
|
||||
import org.springframework.util.ObjectUtils;
|
||||
import org.springframework.util.ReflectionUtils;
|
||||
import org.springframework.util.ReflectionUtils.MethodFilter;
|
||||
import org.springframework.util.StringUtils;
|
||||
|
||||
/**
|
||||
@@ -181,13 +179,13 @@ public class MessagingMethodInvokerHelper extends AbstractExpressionEvaluator im
|
||||
|
||||
private final boolean canProcessMessageList;
|
||||
|
||||
private HandlerMethod handlerMethod;
|
||||
private final String methodName;
|
||||
|
||||
private Class<? extends Annotation> annotationType;
|
||||
private final Method method;
|
||||
|
||||
private String methodName;
|
||||
private final Class<? extends Annotation> annotationType;
|
||||
|
||||
private Method method;
|
||||
private final HandlerMethod handlerMethod;
|
||||
|
||||
private HandlerMethod defaultHandlerMethod;
|
||||
|
||||
@@ -242,6 +240,7 @@ public class MessagingMethodInvokerHelper extends AbstractExpressionEvaluator im
|
||||
this.canProcessMessageList = canProcessMessageList;
|
||||
Assert.notNull(method, "method must not be null");
|
||||
this.method = method;
|
||||
this.methodName = null;
|
||||
this.requiresReply = expectedType != null;
|
||||
if (expectedType != null) {
|
||||
Assert.isTrue(method.getReturnType() != Void.class && method.getReturnType() != Void.TYPE,
|
||||
@@ -260,24 +259,33 @@ public class MessagingMethodInvokerHelper extends AbstractExpressionEvaluator im
|
||||
this.handlerMethodsList.add(
|
||||
Collections.singletonMap(this.handlerMethod.targetParameterType, this.handlerMethod));
|
||||
setDisplayString(targetObject, method);
|
||||
|
||||
JsonObjectMapper<?, ?> mapper;
|
||||
try {
|
||||
mapper = JsonObjectMapperProvider.newInstance();
|
||||
}
|
||||
catch (IllegalStateException e) {
|
||||
mapper = null;
|
||||
}
|
||||
this.jsonObjectMapper = mapper;
|
||||
this.jsonObjectMapper = configureJsonObjectMapperIfAny();
|
||||
}
|
||||
|
||||
private MessagingMethodInvokerHelper(Object targetObject, Class<? extends Annotation> annotationType,
|
||||
String methodName, Class<?> expectedType, boolean canProcessMessageList) {
|
||||
|
||||
this.annotationType = annotationType;
|
||||
this.methodName = methodName;
|
||||
this.canProcessMessageList = canProcessMessageList;
|
||||
Assert.notNull(targetObject, "targetObject must not be null");
|
||||
this.annotationType = annotationType;
|
||||
if (methodName == null) {
|
||||
if (targetObject instanceof Function) {
|
||||
this.methodName = "apply";
|
||||
}
|
||||
else if (targetObject instanceof Consumer) {
|
||||
this.methodName = "accept";
|
||||
}
|
||||
else {
|
||||
this.methodName = null;
|
||||
}
|
||||
}
|
||||
else {
|
||||
this.methodName = methodName;
|
||||
}
|
||||
|
||||
this.method = null;
|
||||
|
||||
this.canProcessMessageList = canProcessMessageList;
|
||||
this.requiresReply = expectedType != null;
|
||||
if (expectedType != null) {
|
||||
this.expectedType = TypeDescriptor.valueOf(expectedType);
|
||||
}
|
||||
@@ -285,8 +293,7 @@ public class MessagingMethodInvokerHelper extends AbstractExpressionEvaluator im
|
||||
this.expectedType = null;
|
||||
}
|
||||
this.targetObject = targetObject;
|
||||
Map<String, Map<Class<?>, HandlerMethod>> handlerMethodsForTarget =
|
||||
findHandlerMethodsForTarget(annotationType, methodName, expectedType != null);
|
||||
Map<String, Map<Class<?>, HandlerMethod>> handlerMethodsForTarget = findHandlerMethodsForTarget();
|
||||
Map<Class<?>, HandlerMethod> methods = handlerMethodsForTarget.get(CANDIDATE_METHODS);
|
||||
Map<Class<?>, HandlerMethod> messageMethods = handlerMethodsForTarget.get(CANDIDATE_MESSAGE_METHODS);
|
||||
if ((methods.size() == 1 && messageMethods.isEmpty()) ||
|
||||
@@ -309,14 +316,16 @@ public class MessagingMethodInvokerHelper extends AbstractExpressionEvaluator im
|
||||
this.handlerMethodsList.add(this.handlerMessageMethods);
|
||||
|
||||
setDisplayString(targetObject, methodName);
|
||||
JsonObjectMapper<?, ?> mapper;
|
||||
this.jsonObjectMapper = configureJsonObjectMapperIfAny();
|
||||
}
|
||||
|
||||
private JsonObjectMapper<?, ?> configureJsonObjectMapperIfAny() {
|
||||
try {
|
||||
mapper = JsonObjectMapperProvider.newInstance();
|
||||
return JsonObjectMapperProvider.newInstance();
|
||||
}
|
||||
catch (IllegalStateException e) {
|
||||
mapper = null;
|
||||
return null;
|
||||
}
|
||||
this.jsonObjectMapper = mapper;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -402,9 +411,9 @@ public class MessagingMethodInvokerHelper extends AbstractExpressionEvaluator im
|
||||
private HandlerMethod createHandlerMethod(Method method) {
|
||||
try {
|
||||
InvocableHandlerMethod invocableHandlerMethod = createInvocableHandlerMethod(method);
|
||||
HandlerMethod handlerMethod = new HandlerMethod(invocableHandlerMethod, this.canProcessMessageList);
|
||||
checkSpelInvokerRequired(getTargetClass(this.targetObject), method, handlerMethod);
|
||||
return handlerMethod;
|
||||
HandlerMethod newHandlerMethod = new HandlerMethod(invocableHandlerMethod, this.canProcessMessageList);
|
||||
checkSpelInvokerRequired(getTargetClass(this.targetObject), method, newHandlerMethod);
|
||||
return newHandlerMethod;
|
||||
}
|
||||
catch (IneligibleMethodException e) {
|
||||
throw new IllegalArgumentException(e);
|
||||
@@ -443,7 +452,7 @@ public class MessagingMethodInvokerHelper extends AbstractExpressionEvaluator im
|
||||
AnnotatedMethodFilter filter = new AnnotatedMethodFilter(this.annotationType, this.methodName,
|
||||
this.requiresReply);
|
||||
Assert.state(canReturnExpectedType(filter, targetType, context.getTypeConverter()),
|
||||
() -> "Cannot convert to expected type (" + this.expectedType + ") from " + this.method);
|
||||
() -> "Cannot convert to expected type (" + this.expectedType + ") from " + this.methodName);
|
||||
context.registerMethodFilter(targetType, filter);
|
||||
}
|
||||
context.setVariable("target", this.targetObject);
|
||||
@@ -753,172 +762,36 @@ public class MessagingMethodInvokerHelper extends AbstractExpressionEvaluator im
|
||||
return contentType != null && contentType.toString().contains("json");
|
||||
}
|
||||
|
||||
private Map<String, Map<Class<?>, HandlerMethod>> findHandlerMethodsForTarget(
|
||||
final Class<? extends Annotation> annotationType, final String methodNameArg, final boolean requiresReply) {
|
||||
|
||||
private Map<String, Map<Class<?>, HandlerMethod>> findHandlerMethodsForTarget() {
|
||||
Map<String, Map<Class<?>, HandlerMethod>> methods = new HashMap<>();
|
||||
Map<Class<?>, HandlerMethod> candidateMethods = new HashMap<>();
|
||||
Map<Class<?>, HandlerMethod> candidateMessageMethods = new HashMap<>();
|
||||
Map<Class<?>, HandlerMethod> fallbackMethods = new HashMap<>();
|
||||
Map<Class<?>, HandlerMethod> fallbackMessageMethods = new HashMap<>();
|
||||
AtomicReference<Class<?>> ambiguousFallbackType = new AtomicReference<>();
|
||||
AtomicReference<Class<?>> ambiguousFallbackMessageGenericType = new AtomicReference<>();
|
||||
Class<?> targetClass = getTargetClass(this.targetObject);
|
||||
|
||||
final Map<Class<?>, HandlerMethod> candidateMethods = new HashMap<>();
|
||||
final Map<Class<?>, HandlerMethod> candidateMessageMethods = new HashMap<>();
|
||||
final Map<Class<?>, HandlerMethod> fallbackMethods = new HashMap<>();
|
||||
final Map<Class<?>, HandlerMethod> fallbackMessageMethods = new HashMap<>();
|
||||
final AtomicReference<Class<?>> ambiguousFallbackType = new AtomicReference<>();
|
||||
final AtomicReference<Class<?>> ambiguousFallbackMessageGenericType = new AtomicReference<>();
|
||||
final Class<?> targetClass = getTargetClass(this.targetObject);
|
||||
|
||||
final String methodNameToUse;
|
||||
|
||||
if (methodNameArg == null) {
|
||||
if (Function.class.isAssignableFrom(targetClass)) {
|
||||
methodNameToUse = "apply";
|
||||
}
|
||||
else if (Consumer.class.isAssignableFrom(targetClass)) {
|
||||
methodNameToUse = "accept";
|
||||
}
|
||||
else {
|
||||
methodNameToUse = null;
|
||||
}
|
||||
}
|
||||
else {
|
||||
methodNameToUse = methodNameArg;
|
||||
}
|
||||
|
||||
|
||||
MethodFilter methodFilter = new UniqueMethodFilter(targetClass);
|
||||
ReflectionUtils.doWithMethods(targetClass, method1 -> {
|
||||
boolean matchesAnnotation = false;
|
||||
if (method1.isBridge()) {
|
||||
return;
|
||||
}
|
||||
if (isMethodDefinedOnObjectClass(method1)) {
|
||||
return;
|
||||
}
|
||||
if (method1.getDeclaringClass().equals(Proxy.class)) {
|
||||
return;
|
||||
}
|
||||
if (annotationType != null && AnnotationUtils.findAnnotation(method1, annotationType) != null) {
|
||||
matchesAnnotation = true;
|
||||
}
|
||||
else if (!Modifier.isPublic(method1.getModifiers())) {
|
||||
return;
|
||||
}
|
||||
if (requiresReply && void.class.equals(method1.getReturnType())) {
|
||||
return;
|
||||
}
|
||||
if (methodNameToUse != null && !methodNameToUse.equals(method1.getName())) {
|
||||
return;
|
||||
}
|
||||
if (methodNameToUse == null
|
||||
&& ObjectUtils.containsElement(new String[] { "start", "stop", "isRunning" }, method1.getName())) {
|
||||
return;
|
||||
}
|
||||
HandlerMethod handlerMethod1;
|
||||
try {
|
||||
method1 = AopUtils.selectInvocableMethod(method1,
|
||||
org.springframework.util.ClassUtils.getUserClass(this.targetObject));
|
||||
handlerMethod1 = createHandlerMethod(method1);
|
||||
}
|
||||
catch (IneligibleMethodException e) {
|
||||
if (LOGGER.isDebugEnabled()) {
|
||||
LOGGER.debug("Method [" + method1 + "] is not eligible for Message handling "
|
||||
+ e.getMessage() + ".");
|
||||
}
|
||||
return;
|
||||
}
|
||||
catch (Exception e) {
|
||||
if (LOGGER.isDebugEnabled()) {
|
||||
LOGGER.debug("Method [" + method1 + "] is not eligible for Message handling.", e);
|
||||
}
|
||||
return;
|
||||
}
|
||||
if (AnnotationUtils.getAnnotation(method1, Default.class) != null) {
|
||||
Assert.state(this.defaultHandlerMethod == null,
|
||||
() -> "Only one method can be @Default, but there are more for: " + this.targetObject);
|
||||
this.defaultHandlerMethod = handlerMethod1;
|
||||
}
|
||||
Class<?> targetParameterType = handlerMethod1.getTargetParameterType();
|
||||
if (matchesAnnotation || annotationType == null) {
|
||||
if (handlerMethod1.isMessageMethod()) {
|
||||
if (candidateMessageMethods.containsKey(targetParameterType)) {
|
||||
throw new IllegalArgumentException("Found more than one method match for type " +
|
||||
"[Message<" + targetParameterType + ">]");
|
||||
}
|
||||
candidateMessageMethods.put(targetParameterType, handlerMethod1);
|
||||
}
|
||||
else {
|
||||
if (candidateMethods.containsKey(targetParameterType)) {
|
||||
String exceptionMessage = "Found more than one method match for ";
|
||||
if (Void.class.equals(targetParameterType)) {
|
||||
exceptionMessage += "empty parameter for 'payload'";
|
||||
}
|
||||
else {
|
||||
exceptionMessage += "type [" + targetParameterType + "]";
|
||||
}
|
||||
throw new IllegalArgumentException(exceptionMessage);
|
||||
}
|
||||
candidateMethods.put(targetParameterType, handlerMethod1);
|
||||
}
|
||||
}
|
||||
else {
|
||||
if (handlerMethod1.isMessageMethod()) {
|
||||
if (fallbackMessageMethods.containsKey(targetParameterType)) {
|
||||
// we need to check for duplicate type matches,
|
||||
// but only if we end up falling back
|
||||
// and we'll only keep track of the first one
|
||||
ambiguousFallbackMessageGenericType.compareAndSet(null, targetParameterType);
|
||||
}
|
||||
fallbackMessageMethods.put(targetParameterType, handlerMethod1);
|
||||
}
|
||||
else {
|
||||
if (fallbackMethods.containsKey(targetParameterType)) {
|
||||
// we need to check for duplicate type matches,
|
||||
// but only if we end up falling back
|
||||
// and we'll only keep track of the first one
|
||||
ambiguousFallbackType.compareAndSet(null, targetParameterType);
|
||||
}
|
||||
fallbackMethods.put(targetParameterType, handlerMethod1);
|
||||
}
|
||||
}
|
||||
}, methodFilter);
|
||||
|
||||
if (candidateMethods.isEmpty() && candidateMessageMethods.isEmpty() && fallbackMethods.isEmpty()
|
||||
&& fallbackMessageMethods.isEmpty()) {
|
||||
findSingleSpecifMethodOnInterfacesIfProxy(methodNameToUse, candidateMessageMethods, candidateMethods);
|
||||
}
|
||||
processMethodsFromTarget(candidateMethods, candidateMessageMethods, fallbackMethods, fallbackMessageMethods,
|
||||
ambiguousFallbackType, ambiguousFallbackMessageGenericType, targetClass);
|
||||
|
||||
if (!candidateMethods.isEmpty() || !candidateMessageMethods.isEmpty()) {
|
||||
methods.put(CANDIDATE_METHODS, candidateMethods);
|
||||
methods.put(CANDIDATE_MESSAGE_METHODS, candidateMessageMethods);
|
||||
return methods;
|
||||
}
|
||||
|
||||
if ((ambiguousFallbackType.get() != null
|
||||
|| ambiguousFallbackMessageGenericType.get() != null)
|
||||
&& ServiceActivator.class.equals(annotationType)) {
|
||||
&& ServiceActivator.class.equals(this.annotationType)) {
|
||||
/*
|
||||
* When there are ambiguous fallback methods,
|
||||
* a Service Activator can finally fallback to RequestReplyExchanger.exchange(m).
|
||||
* Ambiguous means > 1 method that takes the same payload type, or > 1 method
|
||||
* that takes a Message with the same generic type.
|
||||
*/
|
||||
List<Method> frameworkMethods = new ArrayList<>();
|
||||
Class<?>[] allInterfaces = org.springframework.util.ClassUtils.getAllInterfacesForClass(targetClass);
|
||||
for (Class<?> iface : allInterfaces) {
|
||||
try {
|
||||
if ("org.springframework.integration.gateway.RequestReplyExchanger".equals(iface.getName())) {
|
||||
frameworkMethods.add(targetClass.getMethod("exchange", Message.class));
|
||||
if (LOGGER.isDebugEnabled()) {
|
||||
LOGGER.debug(this.targetObject.getClass() +
|
||||
": Ambiguous fallback methods; using RequestReplyExchanger.exchange()");
|
||||
}
|
||||
}
|
||||
}
|
||||
catch (Exception e) {
|
||||
// should never happen (but would fall through to errors below)
|
||||
}
|
||||
}
|
||||
if (frameworkMethods.size() == 1) {
|
||||
Method frameworkMethod = org.springframework.util.ClassUtils.getMostSpecificMethod(
|
||||
frameworkMethods.get(0), this.targetObject.getClass());
|
||||
Method frameworkMethod = obtainFrameworkMethod(targetClass);
|
||||
if (frameworkMethod != null) {
|
||||
HandlerMethod theHandlerMethod = createHandlerMethod(frameworkMethod);
|
||||
methods.put(CANDIDATE_METHODS, Collections.singletonMap(Object.class, theHandlerMethod));
|
||||
methods.put(CANDIDATE_MESSAGE_METHODS, candidateMessageMethods);
|
||||
@@ -926,6 +799,17 @@ public class MessagingMethodInvokerHelper extends AbstractExpressionEvaluator im
|
||||
}
|
||||
}
|
||||
|
||||
validateFallbackMethods(fallbackMethods, fallbackMessageMethods, ambiguousFallbackType,
|
||||
ambiguousFallbackMessageGenericType);
|
||||
|
||||
methods.put(CANDIDATE_METHODS, fallbackMethods);
|
||||
methods.put(CANDIDATE_MESSAGE_METHODS, fallbackMessageMethods);
|
||||
return methods;
|
||||
}
|
||||
|
||||
private void validateFallbackMethods(Map<Class<?>, HandlerMethod> fallbackMethods,
|
||||
Map<Class<?>, HandlerMethod> fallbackMessageMethods, AtomicReference<Class<?>> ambiguousFallbackType,
|
||||
AtomicReference<Class<?>> ambiguousFallbackMessageGenericType) {
|
||||
Assert.state(!fallbackMethods.isEmpty() || !fallbackMessageMethods.isEmpty(),
|
||||
() -> "Target object of type [" + this.targetObject.getClass() +
|
||||
"] has no eligible methods for handling Messages.");
|
||||
@@ -938,14 +822,151 @@ public class MessagingMethodInvokerHelper extends AbstractExpressionEvaluator im
|
||||
+ ambiguousFallbackMessageGenericType
|
||||
+ "] for method match: "
|
||||
+ fallbackMethods.values());
|
||||
|
||||
methods.put(CANDIDATE_METHODS, fallbackMethods);
|
||||
methods.put(CANDIDATE_MESSAGE_METHODS, fallbackMessageMethods);
|
||||
return methods;
|
||||
}
|
||||
|
||||
private void findSingleSpecifMethodOnInterfacesIfProxy(final String methodName,
|
||||
Map<Class<?>, HandlerMethod> candidateMessageMethods,
|
||||
private void processMethodsFromTarget(Map<Class<?>, HandlerMethod> candidateMethods,
|
||||
Map<Class<?>, HandlerMethod> candidateMessageMethods, Map<Class<?>, HandlerMethod> fallbackMethods,
|
||||
Map<Class<?>, HandlerMethod> fallbackMessageMethods, AtomicReference<Class<?>> ambiguousFallbackType,
|
||||
AtomicReference<Class<?>> ambiguousFallbackMessageGenericType, Class<?> targetClass) {
|
||||
|
||||
ReflectionUtils.doWithMethods(targetClass, method1 -> {
|
||||
boolean matchesAnnotation = false;
|
||||
if (this.annotationType != null && AnnotationUtils.findAnnotation(method1, this.annotationType) != null) {
|
||||
matchesAnnotation = true;
|
||||
}
|
||||
else if (!Modifier.isPublic(method1.getModifiers())) {
|
||||
return;
|
||||
}
|
||||
|
||||
HandlerMethod handlerMethod1 = obtainHandlerMethodIfAny(method1);
|
||||
|
||||
if (handlerMethod1 != null) {
|
||||
populateHandlerMethod(candidateMethods, candidateMessageMethods, fallbackMethods,
|
||||
fallbackMessageMethods,
|
||||
ambiguousFallbackType, ambiguousFallbackMessageGenericType, matchesAnnotation, handlerMethod1);
|
||||
}
|
||||
|
||||
}, new UniqueMethodFilter(targetClass));
|
||||
|
||||
if (candidateMethods.isEmpty() && candidateMessageMethods.isEmpty() && fallbackMethods.isEmpty()
|
||||
&& fallbackMessageMethods.isEmpty()) {
|
||||
findSingleSpecifMethodOnInterfacesIfProxy(candidateMessageMethods, candidateMethods);
|
||||
}
|
||||
}
|
||||
|
||||
@Nullable
|
||||
private HandlerMethod obtainHandlerMethodIfAny(Method methodToProcess) {
|
||||
HandlerMethod handlerMethodToUse = null;
|
||||
if (isMethodEligible(methodToProcess)) {
|
||||
try {
|
||||
handlerMethodToUse = createHandlerMethod(
|
||||
AopUtils.selectInvocableMethod(methodToProcess, ClassUtils.getUserClass(this.targetObject)));
|
||||
}
|
||||
catch (Exception e) {
|
||||
if (LOGGER.isDebugEnabled()) {
|
||||
LOGGER.debug("Method [" + methodToProcess + "] is not eligible for Message handling.", e);
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
if (AnnotationUtils.getAnnotation(methodToProcess, Default.class) != null) {
|
||||
Assert.state(this.defaultHandlerMethod == null,
|
||||
() -> "Only one method can be @Default, but there are more for: " + this.targetObject);
|
||||
this.defaultHandlerMethod = handlerMethodToUse;
|
||||
}
|
||||
}
|
||||
|
||||
return handlerMethodToUse;
|
||||
}
|
||||
|
||||
private boolean isMethodEligible(Method methodToProcess) {
|
||||
return !(methodToProcess.isBridge() || // NOSONAR boolean complexity
|
||||
isMethodDefinedOnObjectClass(methodToProcess) ||
|
||||
methodToProcess.getDeclaringClass().equals(Proxy.class) ||
|
||||
(this.requiresReply && void.class.equals(methodToProcess.getReturnType())) ||
|
||||
(this.methodName != null && !this.methodName.equals(methodToProcess.getName())) ||
|
||||
(this.methodName == null &&
|
||||
ObjectUtils.containsElement(new String[] { "start", "stop", "isRunning" },
|
||||
methodToProcess.getName())));
|
||||
}
|
||||
|
||||
private void populateHandlerMethod(Map<Class<?>, HandlerMethod> candidateMethods,
|
||||
Map<Class<?>, HandlerMethod> candidateMessageMethods, Map<Class<?>, HandlerMethod> fallbackMethods,
|
||||
Map<Class<?>, HandlerMethod> fallbackMessageMethods, AtomicReference<Class<?>> ambiguousFallbackType,
|
||||
AtomicReference<Class<?>> ambiguousFallbackMessageGenericType, boolean matchesAnnotation,
|
||||
HandlerMethod handlerMethod1) {
|
||||
|
||||
Class<?> targetParameterType = handlerMethod1.getTargetParameterType();
|
||||
if (matchesAnnotation || this.annotationType == null) {
|
||||
if (handlerMethod1.isMessageMethod()) {
|
||||
if (candidateMessageMethods.containsKey(targetParameterType)) {
|
||||
throw new IllegalArgumentException("Found more than one method match for type " +
|
||||
"[Message<" + targetParameterType + ">]");
|
||||
}
|
||||
candidateMessageMethods.put(targetParameterType, handlerMethod1);
|
||||
}
|
||||
else {
|
||||
if (candidateMethods.containsKey(targetParameterType)) {
|
||||
String exceptionMessage = "Found more than one method match for ";
|
||||
if (Void.class.equals(targetParameterType)) {
|
||||
exceptionMessage += "empty parameter for 'payload'";
|
||||
}
|
||||
else {
|
||||
exceptionMessage += "type [" + targetParameterType + "]";
|
||||
}
|
||||
throw new IllegalArgumentException(exceptionMessage);
|
||||
}
|
||||
candidateMethods.put(targetParameterType, handlerMethod1);
|
||||
}
|
||||
}
|
||||
else {
|
||||
if (handlerMethod1.isMessageMethod()) {
|
||||
if (fallbackMessageMethods.containsKey(targetParameterType)) {
|
||||
// we need to check for duplicate type matches,
|
||||
// but only if we end up falling back
|
||||
// and we'll only keep track of the first one
|
||||
ambiguousFallbackMessageGenericType.compareAndSet(null, targetParameterType);
|
||||
}
|
||||
fallbackMessageMethods.put(targetParameterType, handlerMethod1);
|
||||
}
|
||||
else {
|
||||
if (fallbackMethods.containsKey(targetParameterType)) {
|
||||
// we need to check for duplicate type matches,
|
||||
// but only if we end up falling back
|
||||
// and we'll only keep track of the first one
|
||||
ambiguousFallbackType.compareAndSet(null, targetParameterType);
|
||||
}
|
||||
fallbackMethods.put(targetParameterType, handlerMethod1);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Nullable
|
||||
private Method obtainFrameworkMethod(Class<?> targetClass) {
|
||||
Method frameworkMethod = null;
|
||||
for (Class<?> iface : ClassUtils.getAllInterfacesForClass(targetClass)) {
|
||||
try {
|
||||
// Can't use real class because of package tangle
|
||||
if ("org.springframework.integration.gateway.RequestReplyExchanger".equals(iface.getName())) {
|
||||
frameworkMethod =
|
||||
ClassUtils.getMostSpecificMethod(
|
||||
targetClass.getMethod("exchange", Message.class),
|
||||
this.targetObject.getClass());
|
||||
if (LOGGER.isDebugEnabled()) {
|
||||
LOGGER.debug(this.targetObject.getClass() +
|
||||
": Ambiguous fallback methods; using RequestReplyExchanger.exchange()");
|
||||
}
|
||||
break;
|
||||
}
|
||||
}
|
||||
catch (Exception ex) {
|
||||
throw new IllegalStateException(ex);
|
||||
}
|
||||
}
|
||||
return frameworkMethod;
|
||||
}
|
||||
|
||||
private void findSingleSpecifMethodOnInterfacesIfProxy(Map<Class<?>, HandlerMethod> candidateMessageMethods,
|
||||
Map<Class<?>, HandlerMethod> candidateMethods) {
|
||||
if (AopUtils.isAopProxy(this.targetObject)) {
|
||||
final AtomicReference<Method> targetMethod = new AtomicReference<>();
|
||||
@@ -954,18 +975,18 @@ public class MessagingMethodInvokerHelper extends AbstractExpressionEvaluator im
|
||||
for (Class<?> clazz : interfaces) {
|
||||
ReflectionUtils.doWithMethods(clazz, method1 -> {
|
||||
if (targetMethod.get() != null) {
|
||||
throw new IllegalStateException("Ambiguous method " + methodName + " on " + this.targetObject);
|
||||
throw new IllegalStateException(
|
||||
"Ambiguous method " + this.methodName + " on " + this.targetObject);
|
||||
}
|
||||
else {
|
||||
targetMethod.set(method1);
|
||||
targetClass.set(clazz);
|
||||
}
|
||||
}, method12 -> method12.getName().equals(methodName));
|
||||
}, method12 -> method12.getName().equals(this.methodName));
|
||||
}
|
||||
Method theMethod = targetMethod.get();
|
||||
if (theMethod != null) {
|
||||
theMethod = org.springframework.util.ClassUtils
|
||||
.getMostSpecificMethod(theMethod, this.targetObject.getClass());
|
||||
theMethod = ClassUtils.getMostSpecificMethod(theMethod, this.targetObject.getClass());
|
||||
HandlerMethod theHandlerMethod = createHandlerMethod(theMethod);
|
||||
Class<?> targetParameterType = theHandlerMethod.getTargetParameterType();
|
||||
if (theHandlerMethod.isMessageMethod()) {
|
||||
@@ -1042,8 +1063,7 @@ public class MessagingMethodInvokerHelper extends AbstractExpressionEvaluator im
|
||||
}
|
||||
}
|
||||
}
|
||||
else if (org.springframework.util.ClassUtils.isCglibProxyClass(targetClass)
|
||||
|| targetClass.getSimpleName().contains("$MockitoMock$")) {
|
||||
else if (ClassUtils.isCglibProxyClass(targetClass) || targetClass.getSimpleName().contains("$MockitoMock$")) {
|
||||
Class<?> superClass = targetObject.getClass().getSuperclass();
|
||||
if (!Object.class.equals(superClass)) {
|
||||
targetClass = superClass;
|
||||
@@ -1078,7 +1098,7 @@ public class MessagingMethodInvokerHelper extends AbstractExpressionEvaluator im
|
||||
Set<Class<?>> candidates = methods.keySet();
|
||||
Class<?> match = null;
|
||||
if (!CollectionUtils.isEmpty(candidates)) {
|
||||
match = ClassUtils.findClosestMatch(payloadType, candidates, true);
|
||||
match = org.springframework.integration.util.ClassUtils.findClosestMatch(payloadType, candidates, true);
|
||||
}
|
||||
if (match != null) {
|
||||
return methods.get(match);
|
||||
@@ -1182,113 +1202,12 @@ public class MessagingMethodInvokerHelper extends AbstractExpressionEvaluator im
|
||||
MethodParameter methodParameter = new MethodParameter(method, i);
|
||||
TypeDescriptor parameterTypeDescriptor = new TypeDescriptor(methodParameter);
|
||||
Class<?> parameterType = parameterTypeDescriptor.getObjectType();
|
||||
Type genericParameterType = method.getGenericParameterTypes()[i];
|
||||
Annotation mappingAnnotation =
|
||||
MessagingAnnotationUtils.findMessagePartAnnotation(parameterAnnotations[i], true);
|
||||
if (mappingAnnotation != null) {
|
||||
Class<? extends Annotation> annotationType = mappingAnnotation.annotationType();
|
||||
if (annotationType.equals(Payload.class)) {
|
||||
sb.append("payload");
|
||||
String qualifierExpression = (String) AnnotationUtils.getValue(mappingAnnotation);
|
||||
if (StringUtils.hasText(qualifierExpression)) {
|
||||
sb.append(".")
|
||||
.append(qualifierExpression);
|
||||
}
|
||||
if (!StringUtils.hasText(qualifierExpression)) {
|
||||
this.setExclusiveTargetParameterType(parameterTypeDescriptor, methodParameter);
|
||||
}
|
||||
}
|
||||
if (annotationType.equals(Payloads.class)) {
|
||||
Assert.isTrue(this.canProcessMessageList,
|
||||
"The @Payloads annotation can only be applied " +
|
||||
"if method handler canProcessMessageList.");
|
||||
Assert.isTrue(Collection.class.isAssignableFrom(parameterType),
|
||||
"The @Payloads annotation can only be applied to a Collection-typed parameter.");
|
||||
sb.append("messages.![payload");
|
||||
String qualifierExpression = ((Payloads) mappingAnnotation).value();
|
||||
if (StringUtils.hasText(qualifierExpression)) {
|
||||
sb.append(".")
|
||||
.append(qualifierExpression);
|
||||
}
|
||||
sb.append("]");
|
||||
if (!StringUtils.hasText(qualifierExpression)) {
|
||||
this.setExclusiveTargetParameterType(parameterTypeDescriptor, methodParameter);
|
||||
}
|
||||
}
|
||||
else if (annotationType.equals(Headers.class)) {
|
||||
Assert.isTrue(Map.class.isAssignableFrom(parameterType),
|
||||
"The @Headers annotation can only be applied to a Map-typed parameter.");
|
||||
sb.append("headers");
|
||||
}
|
||||
else if (annotationType.equals(Header.class)) {
|
||||
sb.append(this.determineHeaderExpression(mappingAnnotation, methodParameter));
|
||||
}
|
||||
}
|
||||
else if (parameterTypeDescriptor.isAssignableTo(MESSAGE_TYPE_DESCRIPTOR)) {
|
||||
this.messageMethod = true;
|
||||
sb.append("message");
|
||||
this.setExclusiveTargetParameterType(parameterTypeDescriptor, methodParameter);
|
||||
}
|
||||
else if (this.canProcessMessageList &&
|
||||
(parameterTypeDescriptor.isAssignableTo(MESSAGE_LIST_TYPE_DESCRIPTOR)
|
||||
|| parameterTypeDescriptor.isAssignableTo(MESSAGE_ARRAY_TYPE_DESCRIPTOR))) {
|
||||
sb.append("messages");
|
||||
this.setExclusiveTargetParameterType(parameterTypeDescriptor, methodParameter);
|
||||
}
|
||||
else if (Collection.class.isAssignableFrom(parameterType) || parameterType.isArray()) {
|
||||
if (this.canProcessMessageList) {
|
||||
sb.append("messages.![payload]");
|
||||
}
|
||||
else {
|
||||
sb.append("payload");
|
||||
}
|
||||
this.setExclusiveTargetParameterType(parameterTypeDescriptor, methodParameter);
|
||||
}
|
||||
else if (Iterator.class.isAssignableFrom(parameterType)) {
|
||||
if (this.canProcessMessageList) {
|
||||
Type type = method.getGenericParameterTypes()[i];
|
||||
Type parameterizedType = null;
|
||||
if (type instanceof ParameterizedType) {
|
||||
parameterizedType = ((ParameterizedType) type).getActualTypeArguments()[0];
|
||||
if (parameterizedType instanceof ParameterizedType) {
|
||||
parameterizedType = ((ParameterizedType) parameterizedType).getRawType();
|
||||
}
|
||||
}
|
||||
if (parameterizedType != null && Message.class.isAssignableFrom((Class<?>) parameterizedType)) {
|
||||
sb.append("messages.iterator()");
|
||||
}
|
||||
else {
|
||||
sb.append("messages.![payload].iterator()");
|
||||
}
|
||||
}
|
||||
else {
|
||||
sb.append("payload.iterator()");
|
||||
}
|
||||
this.setExclusiveTargetParameterType(parameterTypeDescriptor, methodParameter);
|
||||
}
|
||||
else if (Map.class.isAssignableFrom(parameterType)) {
|
||||
if (Properties.class.isAssignableFrom(parameterType)) {
|
||||
sb.append("payload instanceof T(java.util.Map) or "
|
||||
+ "(payload instanceof T(String) and payload.contains('=')) ? payload : headers");
|
||||
}
|
||||
else {
|
||||
sb.append("(payload instanceof T(java.util.Map) ? payload : headers)");
|
||||
}
|
||||
Assert.isTrue(!hasUnqualifiedMapParameter,
|
||||
"Found more than one Map typed parameter without any qualification. "
|
||||
+ "Consider using @Payload or @Headers on at least one of the parameters.");
|
||||
hasUnqualifiedMapParameter = true;
|
||||
}
|
||||
else {
|
||||
sb.append("payload");
|
||||
this.setExclusiveTargetParameterType(parameterTypeDescriptor, methodParameter);
|
||||
}
|
||||
}
|
||||
if (hasUnqualifiedMapParameter) {
|
||||
if (this.targetParameterType != null && Map.class.isAssignableFrom(this.targetParameterType)) {
|
||||
throw new IllegalArgumentException(
|
||||
"Unable to determine payload matching parameter due to ambiguous Map typed parameters. "
|
||||
+ "Consider adding the @Payload and or @Headers annotations as appropriate.");
|
||||
}
|
||||
hasUnqualifiedMapParameter = processMethodParameterForExpression(sb, hasUnqualifiedMapParameter,
|
||||
methodParameter, parameterTypeDescriptor, parameterType, genericParameterType,
|
||||
mappingAnnotation);
|
||||
}
|
||||
sb.append(")");
|
||||
if (this.targetParameterTypeDescriptor == null) {
|
||||
@@ -1297,6 +1216,134 @@ public class MessagingMethodInvokerHelper extends AbstractExpressionEvaluator im
|
||||
return sb.toString();
|
||||
}
|
||||
|
||||
private boolean processMethodParameterForExpression(StringBuilder sb, boolean hasUnqualifiedMapParameter,
|
||||
MethodParameter methodParameter, TypeDescriptor parameterTypeDescriptor, Class<?> parameterType,
|
||||
Type genericParameterType, Annotation mappingAnnotation) {
|
||||
|
||||
if (mappingAnnotation != null) {
|
||||
processMappingAnnotationForExpression(sb, methodParameter, parameterTypeDescriptor, parameterType,
|
||||
mappingAnnotation);
|
||||
}
|
||||
else if (parameterTypeDescriptor.isAssignableTo(MESSAGE_TYPE_DESCRIPTOR)) {
|
||||
this.messageMethod = true;
|
||||
sb.append("message");
|
||||
setExclusiveTargetParameterType(parameterTypeDescriptor, methodParameter);
|
||||
}
|
||||
else if (this.canProcessMessageList &&
|
||||
(parameterTypeDescriptor.isAssignableTo(MESSAGE_LIST_TYPE_DESCRIPTOR)
|
||||
|| parameterTypeDescriptor.isAssignableTo(MESSAGE_ARRAY_TYPE_DESCRIPTOR))) {
|
||||
sb.append("messages");
|
||||
setExclusiveTargetParameterType(parameterTypeDescriptor, methodParameter);
|
||||
}
|
||||
else if (Collection.class.isAssignableFrom(parameterType) || parameterType.isArray()) {
|
||||
addCollectionParameterForExpression(sb);
|
||||
setExclusiveTargetParameterType(parameterTypeDescriptor, methodParameter);
|
||||
}
|
||||
else if (Iterator.class.isAssignableFrom(parameterType)) {
|
||||
populateIteratorParameterForExpression(sb, genericParameterType);
|
||||
setExclusiveTargetParameterType(parameterTypeDescriptor, methodParameter);
|
||||
}
|
||||
else if (Map.class.isAssignableFrom(parameterType)) {
|
||||
Assert.isTrue(!hasUnqualifiedMapParameter,
|
||||
"Found more than one Map typed parameter without any qualification. "
|
||||
+ "Consider using @Payload or @Headers on at least one of the parameters.");
|
||||
populateMapParameterForExpression(sb, parameterType);
|
||||
return true;
|
||||
}
|
||||
else {
|
||||
sb.append("payload");
|
||||
setExclusiveTargetParameterType(parameterTypeDescriptor, methodParameter);
|
||||
}
|
||||
return hasUnqualifiedMapParameter;
|
||||
}
|
||||
|
||||
private void processMappingAnnotationForExpression(StringBuilder sb, MethodParameter methodParameter,
|
||||
TypeDescriptor parameterTypeDescriptor, Class<?> parameterType, Annotation mappingAnnotation) {
|
||||
|
||||
Class<? extends Annotation> annotationType = mappingAnnotation.annotationType();
|
||||
if (annotationType.equals(Payload.class)) {
|
||||
sb.append("payload");
|
||||
String qualifierExpression = (String) AnnotationUtils.getValue(mappingAnnotation);
|
||||
if (StringUtils.hasText(qualifierExpression)) {
|
||||
sb.append(".")
|
||||
.append(qualifierExpression);
|
||||
}
|
||||
if (!StringUtils.hasText(qualifierExpression)) {
|
||||
setExclusiveTargetParameterType(parameterTypeDescriptor, methodParameter);
|
||||
}
|
||||
}
|
||||
if (annotationType.equals(Payloads.class)) {
|
||||
Assert.isTrue(this.canProcessMessageList,
|
||||
"The @Payloads annotation can only be applied " +
|
||||
"if method handler canProcessMessageList.");
|
||||
Assert.isTrue(Collection.class.isAssignableFrom(parameterType),
|
||||
"The @Payloads annotation can only be applied to a Collection-typed parameter.");
|
||||
sb.append("messages.![payload");
|
||||
String qualifierExpression = ((Payloads) mappingAnnotation).value();
|
||||
if (StringUtils.hasText(qualifierExpression)) {
|
||||
sb.append(".")
|
||||
.append(qualifierExpression);
|
||||
}
|
||||
sb.append("]");
|
||||
if (!StringUtils.hasText(qualifierExpression)) {
|
||||
setExclusiveTargetParameterType(parameterTypeDescriptor, methodParameter);
|
||||
}
|
||||
}
|
||||
else if (annotationType.equals(Headers.class)) {
|
||||
Assert.isTrue(Map.class.isAssignableFrom(parameterType),
|
||||
"The @Headers annotation can only be applied to a Map-typed parameter.");
|
||||
sb.append("headers");
|
||||
}
|
||||
else if (annotationType.equals(Header.class)) {
|
||||
sb.append(determineHeaderExpression(mappingAnnotation, methodParameter));
|
||||
}
|
||||
}
|
||||
|
||||
private void addCollectionParameterForExpression(StringBuilder sb) {
|
||||
if (this.canProcessMessageList) {
|
||||
sb.append("messages.![payload]");
|
||||
}
|
||||
else {
|
||||
sb.append("payload");
|
||||
}
|
||||
}
|
||||
|
||||
private void populateIteratorParameterForExpression(StringBuilder sb, Type type) {
|
||||
if (this.canProcessMessageList) {
|
||||
Type parameterizedType = null;
|
||||
if (type instanceof ParameterizedType) {
|
||||
parameterizedType = ((ParameterizedType) type).getActualTypeArguments()[0];
|
||||
if (parameterizedType instanceof ParameterizedType) {
|
||||
parameterizedType = ((ParameterizedType) parameterizedType).getRawType();
|
||||
}
|
||||
}
|
||||
if (parameterizedType != null && Message.class.isAssignableFrom((Class<?>) parameterizedType)) {
|
||||
sb.append("messages.iterator()");
|
||||
}
|
||||
else {
|
||||
sb.append("messages.![payload].iterator()");
|
||||
}
|
||||
}
|
||||
else {
|
||||
sb.append("payload.iterator()");
|
||||
}
|
||||
}
|
||||
|
||||
private void populateMapParameterForExpression(StringBuilder sb, Class<?> parameterType) {
|
||||
if (Properties.class.isAssignableFrom(parameterType)) {
|
||||
sb.append("payload instanceof T(java.util.Map) or "
|
||||
+ "(payload instanceof T(String) and payload.contains('=')) ? payload : headers");
|
||||
}
|
||||
else {
|
||||
sb.append("(payload instanceof T(java.util.Map) ? payload : headers)");
|
||||
}
|
||||
if (this.targetParameterType != null && Map.class.isAssignableFrom(this.targetParameterType)) {
|
||||
throw new IllegalArgumentException(
|
||||
"Unable to determine payload matching parameter due to ambiguous Map typed parameters. "
|
||||
+ "Consider adding the @Payload and or @Headers annotations as appropriate.");
|
||||
}
|
||||
}
|
||||
|
||||
private String determineHeaderExpression(Annotation headerAnnotation, MethodParameter methodParameter) {
|
||||
methodParameter.initParameterNameDiscovery(PARAMETER_NAME_DISCOVERER);
|
||||
String headerName = null;
|
||||
|
||||
@@ -54,26 +54,26 @@ public class ZookeeperMetadataStore implements ListenableMetadataStore, SmartLif
|
||||
|
||||
private final CuratorFramework client;
|
||||
|
||||
private final List<MetadataStoreListener> listeners = new CopyOnWriteArrayList<MetadataStoreListener>();
|
||||
private final List<MetadataStoreListener> listeners = new CopyOnWriteArrayList<>();
|
||||
|
||||
/**
|
||||
* An internal map storing local updates, ensuring that they have precedence if the cache contains stale data.
|
||||
* As changes are propagated back from Zookeeper to the cache, entries are removed.
|
||||
*/
|
||||
private final ConcurrentMap<String, LocalChildData> updateMap = new ConcurrentHashMap<String, LocalChildData>();
|
||||
private final ConcurrentMap<String, LocalChildData> updateMap = new ConcurrentHashMap<>();
|
||||
|
||||
private volatile String root = "/SpringIntegration-MetadataStore";
|
||||
private String root = "/SpringIntegration-MetadataStore";
|
||||
|
||||
private volatile String encoding = "UTF-8";
|
||||
private String encoding = "UTF-8";
|
||||
|
||||
private volatile PathChildrenCache cache;
|
||||
private PathChildrenCache cache;
|
||||
|
||||
private boolean autoStartup = true;
|
||||
|
||||
private int phase = Integer.MAX_VALUE;
|
||||
|
||||
private volatile boolean running = false;
|
||||
|
||||
private volatile boolean autoStartup = true;
|
||||
|
||||
private volatile int phase = Integer.MAX_VALUE;
|
||||
|
||||
public ZookeeperMetadataStore(CuratorFramework client) {
|
||||
Assert.notNull(client, "Client cannot be null");
|
||||
this.client = client;
|
||||
@@ -81,7 +81,6 @@ public class ZookeeperMetadataStore implements ListenableMetadataStore, SmartLif
|
||||
|
||||
/**
|
||||
* Encoding to use when storing data in ZooKeeper
|
||||
*
|
||||
* @param encoding encoding as text
|
||||
*/
|
||||
public void setEncoding(String encoding) {
|
||||
@@ -91,7 +90,6 @@ public class ZookeeperMetadataStore implements ListenableMetadataStore, SmartLif
|
||||
|
||||
/**
|
||||
* Root node - store entries are children of this node.
|
||||
*
|
||||
* @param root encoding as text
|
||||
*/
|
||||
public void setRoot(String root) {
|
||||
@@ -124,14 +122,7 @@ public class ZookeeperMetadataStore implements ListenableMetadataStore, SmartLif
|
||||
}
|
||||
catch (@SuppressWarnings(UNUSED) KeeperException.NodeExistsException e) {
|
||||
// so the data actually exists, we can read it
|
||||
try {
|
||||
byte[] bytes = this.client.getData().forPath(getPath(key));
|
||||
return IntegrationUtils.bytesToString(bytes, this.encoding);
|
||||
}
|
||||
catch (Exception exceptionDuringGet) {
|
||||
throw new ZookeeperMetadataStoreException("Exception while reading node with key '" + key + "':",
|
||||
exceptionDuringGet);
|
||||
}
|
||||
return get(key);
|
||||
}
|
||||
catch (Exception e) {
|
||||
throw new ZookeeperMetadataStoreException("Error while trying to set '" + key + "':", e);
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2015-2016 the original author or authors.
|
||||
* Copyright 2015-2019 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.
|
||||
@@ -33,6 +33,8 @@ import org.junit.BeforeClass;
|
||||
/**
|
||||
* @author Marius Bogoevici
|
||||
* @author Gary Russell
|
||||
* @author Artem Bilan
|
||||
*
|
||||
* @since 4.2
|
||||
*
|
||||
*/
|
||||
@@ -52,7 +54,7 @@ public class ZookeeperTestSupport {
|
||||
}
|
||||
|
||||
@AfterClass
|
||||
public static void tearDownClass() throws Exception {
|
||||
public static void tearDownClass() {
|
||||
try {
|
||||
testingServer.stop();
|
||||
}
|
||||
@@ -63,7 +65,7 @@ public class ZookeeperTestSupport {
|
||||
}
|
||||
|
||||
@Before
|
||||
public void setUp() throws Exception {
|
||||
public void setUp() {
|
||||
client = createNewClient();
|
||||
}
|
||||
|
||||
@@ -72,7 +74,7 @@ public class ZookeeperTestSupport {
|
||||
CloseableUtils.closeQuietly(this.client);
|
||||
}
|
||||
|
||||
protected static CuratorFramework createNewClient() throws InterruptedException {
|
||||
protected static CuratorFramework createNewClient() {
|
||||
CuratorFramework client = CuratorFrameworkFactory.newClient(testingServer.getConnectString(),
|
||||
new BoundedExponentialBackoffRetry(100, 1000, 3));
|
||||
client.start();
|
||||
|
||||
@@ -53,7 +53,7 @@ public class ZookeeperMetadataStoreTests extends ZookeeperTestSupport {
|
||||
|
||||
@Override
|
||||
@Before
|
||||
public void setUp() throws Exception {
|
||||
public void setUp() {
|
||||
super.setUp();
|
||||
this.metadataStore = new ZookeeperMetadataStore(client);
|
||||
this.metadataStore.start();
|
||||
@@ -83,7 +83,7 @@ public class ZookeeperMetadataStoreTests extends ZookeeperTestSupport {
|
||||
|
||||
|
||||
@Test
|
||||
public void testGetValueFromMetadataStore() throws Exception {
|
||||
public void testGetValueFromMetadataStore() {
|
||||
String testKey = "ZookeeperMetadataStoreTests-GetValue";
|
||||
metadataStore.put(testKey, "Hello Zookeeper");
|
||||
String retrievedValue = metadataStore.get(testKey);
|
||||
@@ -194,7 +194,7 @@ public class ZookeeperMetadataStoreTests extends ZookeeperTestSupport {
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testRemoveFromMetadataStore() throws Exception {
|
||||
public void testRemoveFromMetadataStore() {
|
||||
String testKey = "ZookeeperMetadataStoreTests-Remove";
|
||||
String testValue = "Integration";
|
||||
metadataStore.put(testKey, testValue);
|
||||
@@ -203,12 +203,12 @@ public class ZookeeperMetadataStoreTests extends ZookeeperTestSupport {
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testListenerInvokedOnLocalChanges() throws Exception {
|
||||
public void testListenerInvokedOnLocalChanges() {
|
||||
String testKey = "ZookeeperMetadataStoreTests";
|
||||
|
||||
// register listeners
|
||||
final List<List<String>> notifiedChanges = new ArrayList<List<String>>();
|
||||
final Map<String, CyclicBarrier> barriers = new HashMap<String, CyclicBarrier>();
|
||||
final List<List<String>> notifiedChanges = new ArrayList<>();
|
||||
final Map<String, CyclicBarrier> barriers = new HashMap<>();
|
||||
barriers.put("add", new CyclicBarrier(2));
|
||||
barriers.put("remove", new CyclicBarrier(2));
|
||||
barriers.put("update", new CyclicBarrier(2));
|
||||
@@ -275,15 +275,16 @@ public class ZookeeperMetadataStoreTests extends ZookeeperTestSupport {
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testListenerInvokedOnRemoteChanges() throws Exception {
|
||||
public void testListenerInvokedOnRemoteChanges() {
|
||||
String testKey = "ZookeeperMetadataStoreTests";
|
||||
|
||||
CuratorFramework otherClient = createNewClient();
|
||||
ZookeeperMetadataStore otherMetadataStore = new ZookeeperMetadataStore(otherClient);
|
||||
otherMetadataStore.start();
|
||||
|
||||
// register listeners
|
||||
final List<List<String>> notifiedChanges = new ArrayList<List<String>>();
|
||||
final Map<String, CyclicBarrier> barriers = new HashMap<String, CyclicBarrier>();
|
||||
final List<List<String>> notifiedChanges = new ArrayList<>();
|
||||
final Map<String, CyclicBarrier> barriers = new HashMap<>();
|
||||
barriers.put("add", new CyclicBarrier(2));
|
||||
barriers.put("remove", new CyclicBarrier(2));
|
||||
barriers.put("update", new CyclicBarrier(2));
|
||||
@@ -346,7 +347,7 @@ public class ZookeeperMetadataStoreTests extends ZookeeperTestSupport {
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testAddRemoveListener() throws Exception {
|
||||
public void testAddRemoveListener() {
|
||||
MetadataStoreListener mockListener = Mockito.mock(MetadataStoreListener.class);
|
||||
DirectFieldAccessor accessor = new DirectFieldAccessor(metadataStore);
|
||||
|
||||
|
||||
Reference in New Issue
Block a user