diff --git a/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractKeyValueMessageStore.java b/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractKeyValueMessageStore.java index 911a0320ee..a691c06af9 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractKeyValueMessageStore.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractKeyValueMessageStore.java @@ -138,7 +138,7 @@ public abstract class AbstractKeyValueMessageStore extends AbstractMessageGroupS for (Message message : messages) { // enrich Message with additional headers and add it to MS Message enrichedMessage = enrichMessage(message); - doAddMessage(message); + doAddMessage(enrichedMessage); if (metadata != null) { metadata.add(enrichedMessage.getHeaders().getId()); } diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/store/RedisMessageGroupStoreTests.java b/spring-integration-redis/src/test/java/org/springframework/integration/redis/store/RedisMessageGroupStoreTests.java index f88bd242c1..df4a8c7fa0 100644 --- a/spring-integration-redis/src/test/java/org/springframework/integration/redis/store/RedisMessageGroupStoreTests.java +++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/store/RedisMessageGroupStoreTests.java @@ -473,7 +473,7 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests { fail("SerializationException expected"); } catch (Exception e) { - assertThat(e.getCause().getCause(), instanceOf(IllegalArgumentException.class)); + assertThat(e.getCause(), instanceOf(IllegalArgumentException.class)); assertThat(e.getMessage(), containsString("The class with " + "org.springframework.integration.redis.store.RedisMessageGroupStoreTests$Foo and name of " + @@ -489,7 +489,7 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests { store.removeMessageGroup(1); messageGroup = store.addMessageToGroup(1, fooMessage); assertEquals(1, messageGroup.size()); - assertEquals(fooMessage, messageGroup.getMessages().iterator().next()); + assertEquals(fooMessage.getPayload(), messageGroup.getMessages().iterator().next().getPayload()); mapper = JacksonJsonUtils.messagingAwareMapper("*"); @@ -499,7 +499,7 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests { store.removeMessageGroup(1); messageGroup = store.addMessageToGroup(1, fooMessage); assertEquals(1, messageGroup.size()); - assertEquals(fooMessage, messageGroup.getMessages().iterator().next()); + assertEquals(fooMessage.getPayload(), messageGroup.getMessages().iterator().next().getPayload()); } private static class Foo {