diff --git a/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java index 42725cf30d..59c20fbd33 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java @@ -31,6 +31,8 @@ import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; import java.util.stream.Stream; +import com.fasterxml.jackson.annotation.JsonCreator; +import com.fasterxml.jackson.annotation.JsonProperty; import org.aopalliance.aop.Advice; import org.springframework.aop.framework.ProxyFactory; @@ -97,6 +99,7 @@ import org.springframework.util.ObjectUtils; * @author Artem Bilan * @author Gary Russell * @author Christian Tzolov + * @author Youbin Wu * * @since 1.0.3 */ @@ -688,7 +691,10 @@ public class DelayHandler extends AbstractReplyProducingMessageHandler implement private final Message original; - DelayedMessageWrapper(Message original, long requestDate) { + @JsonCreator(mode = JsonCreator.Mode.PROPERTIES) + DelayedMessageWrapper(@JsonProperty("original") Message original, + @JsonProperty("requestDate") long requestDate) { + this.original = original; this.requestDate = requestDate; } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/json/JacksonJsonUtils.java b/spring-integration-core/src/main/java/org/springframework/integration/support/json/JacksonJsonUtils.java index 8beed4b565..f8e7076dc1 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/support/json/JacksonJsonUtils.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/support/json/JacksonJsonUtils.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2022 the original author or authors. + * Copyright 2002-2024 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. @@ -46,6 +46,7 @@ import org.springframework.messaging.support.GenericMessage; * * @author Artem Bilan * @author Gary Russell + * @author Youbin Wu * * @since 3.0 * @@ -63,7 +64,8 @@ public final class JacksonJsonUtils { "org.springframework.integration.support", "org.springframework.integration.message", "org.springframework.integration.store", - "org.springframework.integration.history" + "org.springframework.integration.history", + "org.springframework.integration.handler" ); private JacksonJsonUtils() { 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 1da003ac88..ec33d8b74a 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 @@ -20,7 +20,6 @@ import java.util.ArrayList; import java.util.Date; import java.util.Iterator; import java.util.List; -import java.util.Objects; import java.util.Properties; import java.util.UUID; import java.util.concurrent.ExecutorService; @@ -35,13 +34,16 @@ import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Disabled; import org.junit.jupiter.api.Test; +import org.springframework.beans.BeanUtils; import org.springframework.context.support.ClassPathXmlApplicationContext; import org.springframework.data.redis.connection.RedisConnectionFactory; import org.springframework.data.redis.core.StringRedisTemplate; import org.springframework.data.redis.serializer.GenericJackson2JsonRedisSerializer; +import org.springframework.data.redis.serializer.SerializationException; import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.channel.NullChannel; import org.springframework.integration.channel.QueueChannel; +import org.springframework.integration.handler.DelayHandler; import org.springframework.integration.history.MessageHistory; import org.springframework.integration.message.AdviceMessage; import org.springframework.integration.redis.RedisContainerTest; @@ -56,14 +58,15 @@ import org.springframework.messaging.support.ErrorMessage; import org.springframework.messaging.support.GenericMessage; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatExceptionOfType; import static org.assertj.core.api.Assertions.assertThatNoException; -import static org.assertj.core.api.Assertions.fail; /** * @author Oleg Zhurakousky * @author Artem Bilan * @author Gary Russell * @author Artem Vozhdayenko + * @author Youbin Wu */ class RedisMessageGroupStoreTests implements RedisContainerTest { @@ -317,7 +320,7 @@ class RedisMessageGroupStoreTests implements RedisContainerTest { executor.execute(() -> { store2.removeMessagesFromGroup(this.groupId, message); MessageGroup group = store2.getMessageGroup(this.groupId); - if (group.getMessages().size() != 0) { + if (!group.getMessages().isEmpty()) { failures.add("REMOVE"); throw new AssertionFailedError("Failed on Remove"); } @@ -401,11 +404,17 @@ class RedisMessageGroupStoreTests implements RedisContainerTest { Message mutableMessage = new MutableMessage<>(UUID.randomUUID()); Message adviceMessage = new AdviceMessage<>("foo", genericMessage); ErrorMessage errorMessage = new ErrorMessage(new RuntimeException("test exception"), mutableMessage); + var delayedMessageWrapperConstructor = + BeanUtils.getResolvableConstructor(DelayHandler.DelayedMessageWrapper.class); + Message delayMessage = new GenericMessage<>( + BeanUtils.instantiateClass(delayedMessageWrapperConstructor, genericMessage, + System.currentTimeMillis())); - store.addMessagesToGroup(this.groupId, genericMessage, mutableMessage, adviceMessage, errorMessage); + store.addMessagesToGroup(this.groupId, + genericMessage, mutableMessage, adviceMessage, errorMessage, delayMessage); MessageGroup messageGroup = store.getMessageGroup(this.groupId); - assertThat(messageGroup.size()).isEqualTo(4); + assertThat(messageGroup.size()).isEqualTo(5); List> messages = new ArrayList<>(messageGroup.getMessages()); assertThat(messages.get(0)).isEqualTo(genericMessage); assertThat(messages.get(0).getHeaders()).containsKeys(MessageHistory.HEADER_NAME); @@ -418,22 +427,21 @@ class RedisMessageGroupStoreTests implements RedisContainerTest { .isEqualTo(errorMessage.getOriginalMessage()); assertThat(((ErrorMessage) errorMessageResult).getPayload().getMessage()) .isEqualTo(errorMessage.getPayload().getMessage()); + assertThat(messages.get(4)).isEqualTo(delayMessage); Message fooMessage = new GenericMessage<>(new Foo("foo")); - try { - store.addMessageToGroup(this.groupId, fooMessage) - .getMessages() - .iterator() - .next(); - fail("SerializationException expected"); - } - catch (Exception e) { - assertThat(e.getCause().getCause()).isInstanceOf(IllegalArgumentException.class); - assertThat(e.getMessage()).contains("The class with " + - "org.springframework.integration.redis.store.RedisMessageGroupStoreTests$Foo and name of " + - "org.springframework.integration.redis.store.RedisMessageGroupStoreTests$Foo " + - "is not in the trusted packages:"); - } + + assertThatExceptionOfType(SerializationException.class) + .isThrownBy(() -> + store.addMessageToGroup(this.groupId, fooMessage) + .getMessages() + .iterator() + .next()) + .withRootCauseInstanceOf(IllegalArgumentException.class) + .withMessageContaining("The class with " + + "org.springframework.integration.redis.store.RedisMessageGroupStoreTests$Foo and name of " + + "org.springframework.integration.redis.store.RedisMessageGroupStoreTests$Foo " + + "is not in the trusted packages:"); mapper = JacksonJsonUtils.messagingAwareMapper(getClass().getPackage().getName()); @@ -486,43 +494,7 @@ class RedisMessageGroupStoreTests implements RedisContainerTest { assertThat(store.messageGroupSize("2")).isEqualTo(1); } - private static class Foo { - - private String foo; - - Foo() { - } - - Foo(String foo) { - this.foo = foo; - } - - public String getFoo() { - return this.foo; - } - - public void setFoo(String foo) { - this.foo = foo; - } - - @Override - public boolean equals(Object o) { - if (this == o) { - return true; - } - if (o == null || getClass() != o.getClass()) { - return false; - } - - Foo foo1 = (Foo) o; - - return this.foo != null ? this.foo.equals(foo1.foo) : foo1.foo == null; - } - - @Override - public int hashCode() { - return Objects.hashCode(this.foo); - } + private record Foo(String foo) { }