From 8843b1ade5e47dfa18b07afc70e51b0bda8049b4 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Wed, 23 Mar 2022 17:17:33 -0400 Subject: [PATCH] GH-1441: Fix Payload Detection with MessageHeaders Resolves https://github.com/spring-projects/spring-amqp/issues/1441 Previously, `MessageHeaders` had to be annotated with `@Headers` so that it was ignored during payload parameter resolution; otherwise it caused ambiguity. Ignore `MessageHeaders` even when not so annotated. Also fix some tests that were checking the same topic and count down latch so were unconditionally passing. Change one of those tests to verify the fix. **cherry-pick to 2.4.x, 2.3.x** --- .../MessagingMessageListenerAdapter.java | 6 ++-- .../EnableRabbitIntegrationTests.java | 30 ++++++++++++------- 2 files changed, 23 insertions(+), 13 deletions(-) diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/MessagingMessageListenerAdapter.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/MessagingMessageListenerAdapter.java index fc585a1b..baba04c4 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/MessagingMessageListenerAdapter.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/MessagingMessageListenerAdapter.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2021 the original author or authors. + * Copyright 2002-2022 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,7 @@ import org.springframework.amqp.support.converter.MessagingMessageConverter; import org.springframework.core.MethodParameter; import org.springframework.lang.Nullable; import org.springframework.messaging.Message; +import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.MessagingException; import org.springframework.messaging.handler.annotation.Header; import org.springframework.messaging.handler.annotation.Headers; @@ -379,7 +380,8 @@ public class MessagingMessageListenerAdapter extends AbstractAdaptableMessageLis * We ignore parameters with type Message because they are not involved with conversion. */ boolean isHeaderOrHeaders = methodParameter.hasParameterAnnotation(Header.class) - || methodParameter.hasParameterAnnotation(Headers.class); + || methodParameter.hasParameterAnnotation(Headers.class) + || methodParameter.getParameterType().equals(MessageHeaders.class); boolean isPayload = methodParameter.hasParameterAnnotation(Payload.class); if (isHeaderOrHeaders && isPayload && MessagingMessageListenerAdapter.this.logger.isWarnEnabled()) { MessagingMessageListenerAdapter.this.logger.warn(this.method.getName() diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/annotation/EnableRabbitIntegrationTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/annotation/EnableRabbitIntegrationTests.java index 850cf110..3b2bee26 100644 --- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/annotation/EnableRabbitIntegrationTests.java +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/annotation/EnableRabbitIntegrationTests.java @@ -120,6 +120,7 @@ import org.springframework.core.convert.support.DefaultConversionService; import org.springframework.core.task.TaskExecutor; import org.springframework.data.web.JsonPath; import org.springframework.lang.NonNull; +import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.converter.GenericMessageConverter; import org.springframework.messaging.handler.annotation.Header; import org.springframework.messaging.handler.annotation.Payload; @@ -560,8 +561,9 @@ public class EnableRabbitIntegrationTests { public void testDifferentTypes() throws InterruptedException { Foo1 foo = new Foo1(); foo.setBar("bar"); + this.service.foos.clear(); this.jsonRabbitTemplate.convertAndSend("differentTypes", foo); - assertThat(this.service.latch.await(10, TimeUnit.SECONDS)).isTrue(); + assertThat(this.service.dtLatch1.await(10, TimeUnit.SECONDS)).isTrue(); assertThat(this.service.foos.get(0)).isInstanceOf(Foo2.class); assertThat(((Foo2) this.service.foos.get(0)).getBar()).isEqualTo("bar"); assertThat(TestUtils.getPropertyValue(this.registry.getListenerContainer("different"), "concurrentConsumers")).isEqualTo(2); @@ -571,8 +573,9 @@ public class EnableRabbitIntegrationTests { public void testDifferentTypesWithConcurrency() throws InterruptedException { Foo1 foo = new Foo1(); foo.setBar("bar"); - this.jsonRabbitTemplate.convertAndSend("differentTypes", foo); - assertThat(this.service.latch.await(10, TimeUnit.SECONDS)).isTrue(); + this.service.foos.clear(); + this.jsonRabbitTemplate.convertAndSend("differentTypes2", foo); + assertThat(this.service.dtLatch2.await(10, TimeUnit.SECONDS)).isTrue(); assertThat(this.service.foos.get(0)).isInstanceOf(Foo2.class); assertThat(((Foo2) this.service.foos.get(0)).getBar()).isEqualTo("bar"); MessageListenerContainer container = this.registry.getListenerContainer("differentWithConcurrency"); @@ -584,8 +587,9 @@ public class EnableRabbitIntegrationTests { public void testDifferentTypesWithVariableConcurrency() throws InterruptedException { Foo1 foo = new Foo1(); foo.setBar("bar"); - this.jsonRabbitTemplate.convertAndSend("differentTypes", foo); - assertThat(this.service.latch.await(10, TimeUnit.SECONDS)).isTrue(); + this.service.foos.clear(); + this.jsonRabbitTemplate.convertAndSend("differentTypes3", foo); + assertThat(this.service.dtLatch3.await(10, TimeUnit.SECONDS)).isTrue(); assertThat(this.service.foos.get(0)).isInstanceOf(Foo2.class); assertThat(((Foo2) this.service.foos.get(0)).getBar()).isEqualTo("bar"); MessageListenerContainer container = this.registry.getListenerContainer("differentWithVariableConcurrency"); @@ -1086,7 +1090,11 @@ public class EnableRabbitIntegrationTests { final List foos = new ArrayList<>(); - final CountDownLatch latch = new CountDownLatch(1); + final CountDownLatch dtLatch1 = new CountDownLatch(1); + + final CountDownLatch dtLatch2 = new CountDownLatch(1); + + final CountDownLatch dtLatch3 = new CountDownLatch(1); final CountDownLatch validationLatch = new CountDownLatch(1); @@ -1237,21 +1245,21 @@ public class EnableRabbitIntegrationTests { containerFactory = "jsonListenerContainerFactoryNoClassMapper") public void handleDifferent(@Validated Foo2 foo) { foos.add(foo); - latch.countDown(); + dtLatch1.countDown(); } @RabbitListener(id = "differentWithConcurrency", queues = "differentTypes2", - containerFactory = "jsonListenerContainerFactory", concurrency = "#{3}") - public void handleDifferentWithConcurrency(Foo2 foo) { + containerFactory = "jsonListenerContainerFactoryNoClassMapper", concurrency = "#{3}") + public void handleDifferentWithConcurrency(Foo2 foo, MessageHeaders headers) { foos.add(foo); - latch.countDown(); + dtLatch2.countDown(); } @RabbitListener(id = "differentWithVariableConcurrency", queues = "differentTypes3", containerFactory = "jsonListenerContainerFactory", concurrency = "3-4") public void handleDifferentWithVariableConcurrency(Foo2 foo) { foos.add(foo); - latch.countDown(); + dtLatch3.countDown(); } @RabbitListener(id = "notStarted", containerFactory = "rabbitAutoStartFalseListenerContainerFactory",