diff --git a/spring-integration-core/src/main/java/org/springframework/integration/handler/support/MessagingMethodInvokerHelper.java b/spring-integration-core/src/main/java/org/springframework/integration/handler/support/MessagingMethodInvokerHelper.java index 72ad7d514b..a4893b191b 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/handler/support/MessagingMethodInvokerHelper.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/handler/support/MessagingMethodInvokerHelper.java @@ -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 annotationType; + private final Method method; - private String methodName; + private final Class 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 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, HandlerMethod>> handlerMethodsForTarget = - findHandlerMethodsForTarget(annotationType, methodName, expectedType != null); + Map, HandlerMethod>> handlerMethodsForTarget = findHandlerMethodsForTarget(); Map, HandlerMethod> methods = handlerMethodsForTarget.get(CANDIDATE_METHODS); Map, 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, HandlerMethod>> findHandlerMethodsForTarget( - final Class annotationType, final String methodNameArg, final boolean requiresReply) { - + private Map, HandlerMethod>> findHandlerMethodsForTarget() { Map, HandlerMethod>> methods = new HashMap<>(); + Map, HandlerMethod> candidateMethods = new HashMap<>(); + Map, HandlerMethod> candidateMessageMethods = new HashMap<>(); + Map, HandlerMethod> fallbackMethods = new HashMap<>(); + Map, HandlerMethod> fallbackMessageMethods = new HashMap<>(); + AtomicReference> ambiguousFallbackType = new AtomicReference<>(); + AtomicReference> ambiguousFallbackMessageGenericType = new AtomicReference<>(); + Class targetClass = getTargetClass(this.targetObject); - final Map, HandlerMethod> candidateMethods = new HashMap<>(); - final Map, HandlerMethod> candidateMessageMethods = new HashMap<>(); - final Map, HandlerMethod> fallbackMethods = new HashMap<>(); - final Map, HandlerMethod> fallbackMessageMethods = new HashMap<>(); - final AtomicReference> ambiguousFallbackType = new AtomicReference<>(); - final AtomicReference> 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 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, HandlerMethod> fallbackMethods, + Map, HandlerMethod> fallbackMessageMethods, AtomicReference> ambiguousFallbackType, + AtomicReference> 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, HandlerMethod> candidateMessageMethods, + private void processMethodsFromTarget(Map, HandlerMethod> candidateMethods, + Map, HandlerMethod> candidateMessageMethods, Map, HandlerMethod> fallbackMethods, + Map, HandlerMethod> fallbackMessageMethods, AtomicReference> ambiguousFallbackType, + AtomicReference> 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, HandlerMethod> candidateMethods, + Map, HandlerMethod> candidateMessageMethods, Map, HandlerMethod> fallbackMethods, + Map, HandlerMethod> fallbackMessageMethods, AtomicReference> ambiguousFallbackType, + AtomicReference> 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, HandlerMethod> candidateMessageMethods, Map, HandlerMethod> candidateMethods) { if (AopUtils.isAopProxy(this.targetObject)) { final AtomicReference 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> 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 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 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; diff --git a/spring-integration-zookeeper/src/main/java/org/springframework/integration/zookeeper/metadata/ZookeeperMetadataStore.java b/spring-integration-zookeeper/src/main/java/org/springframework/integration/zookeeper/metadata/ZookeeperMetadataStore.java index b635d02af9..7a19adb57c 100644 --- a/spring-integration-zookeeper/src/main/java/org/springframework/integration/zookeeper/metadata/ZookeeperMetadataStore.java +++ b/spring-integration-zookeeper/src/main/java/org/springframework/integration/zookeeper/metadata/ZookeeperMetadataStore.java @@ -54,26 +54,26 @@ public class ZookeeperMetadataStore implements ListenableMetadataStore, SmartLif private final CuratorFramework client; - private final List listeners = new CopyOnWriteArrayList(); + private final List 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 updateMap = new ConcurrentHashMap(); + private final ConcurrentMap 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); diff --git a/spring-integration-zookeeper/src/test/java/org/springframework/integration/zookeeper/ZookeeperTestSupport.java b/spring-integration-zookeeper/src/test/java/org/springframework/integration/zookeeper/ZookeeperTestSupport.java index 64384724e2..be3e14784d 100644 --- a/spring-integration-zookeeper/src/test/java/org/springframework/integration/zookeeper/ZookeeperTestSupport.java +++ b/spring-integration-zookeeper/src/test/java/org/springframework/integration/zookeeper/ZookeeperTestSupport.java @@ -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(); diff --git a/spring-integration-zookeeper/src/test/java/org/springframework/integration/zookeeper/metadata/ZookeeperMetadataStoreTests.java b/spring-integration-zookeeper/src/test/java/org/springframework/integration/zookeeper/metadata/ZookeeperMetadataStoreTests.java index 91a0b6dd32..fbba3bd6d2 100644 --- a/spring-integration-zookeeper/src/test/java/org/springframework/integration/zookeeper/metadata/ZookeeperMetadataStoreTests.java +++ b/spring-integration-zookeeper/src/test/java/org/springframework/integration/zookeeper/metadata/ZookeeperMetadataStoreTests.java @@ -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> notifiedChanges = new ArrayList>(); - final Map barriers = new HashMap(); + final List> notifiedChanges = new ArrayList<>(); + final Map 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> notifiedChanges = new ArrayList>(); - final Map barriers = new HashMap(); + final List> notifiedChanges = new ArrayList<>(); + final Map 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);