INT-3551: Idempotent Receiver: Add value-strategy
JIRA: https://jira.spring.io/browse/INT-3551 * Rename `MetadataKeyStrategy` -> `MetadataEntryStrategy` * Add `valueStrategy` to the `MetadataStoreSelector` * Add `value-strategy` and `value-expression` to the `<idempotent-receiver>` INT-3551: Add `@IR` support on service methods Rework `MetadataKeyStrategy` just to the `MessageProcessor` Fix Docs Minor Doc Polishing.
This commit is contained in:
committed by
Gary Russell
parent
1018fcf361
commit
0cc9273a2e
@@ -23,6 +23,8 @@ import static org.junit.Assert.assertThat;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
import static org.junit.Assert.fail;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
@@ -36,14 +38,15 @@ import org.springframework.context.annotation.Configuration;
|
||||
import org.springframework.integration.IntegrationMessageHeaderAccessor;
|
||||
import org.springframework.integration.MessageRejectedException;
|
||||
import org.springframework.integration.annotation.IdempotentReceiver;
|
||||
import org.springframework.integration.annotation.ServiceActivator;
|
||||
import org.springframework.integration.channel.DirectChannel;
|
||||
import org.springframework.integration.channel.QueueChannel;
|
||||
import org.springframework.integration.config.EnableIntegration;
|
||||
import org.springframework.integration.handler.MessageProcessor;
|
||||
import org.springframework.integration.handler.advice.AbstractRequestHandlerAdvice;
|
||||
import org.springframework.integration.handler.advice.IdempotentReceiverInterceptor;
|
||||
import org.springframework.integration.jmx.config.EnableIntegrationMBeanExport;
|
||||
import org.springframework.integration.metadata.ConcurrentMetadataStore;
|
||||
import org.springframework.integration.metadata.MetadataKeyStrategy;
|
||||
import org.springframework.integration.metadata.MetadataStore;
|
||||
import org.springframework.integration.metadata.SimpleMetadataStore;
|
||||
import org.springframework.integration.selector.MetadataStoreSelector;
|
||||
@@ -54,6 +57,7 @@ import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.MessageChannel;
|
||||
import org.springframework.messaging.PollableChannel;
|
||||
import org.springframework.messaging.support.GenericMessage;
|
||||
import org.springframework.stereotype.Component;
|
||||
import org.springframework.test.annotation.DirtiesContext;
|
||||
import org.springframework.test.context.ContextConfiguration;
|
||||
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
|
||||
@@ -82,15 +86,24 @@ public class IdempotentReceiverIntegrationTests {
|
||||
@Autowired
|
||||
private AtomicInteger adviceCalled;
|
||||
|
||||
@Autowired
|
||||
private MessageChannel annotatedMethodChannel;
|
||||
|
||||
@Autowired
|
||||
private FooService fooService;
|
||||
|
||||
@Test
|
||||
public void testIdempotentReceiver() {
|
||||
this.idempotentReceiverInterceptor.setThrowExceptionOnRejection(true);
|
||||
TestUtils.getPropertyValue(this.store, "metadata", Map.class).clear();
|
||||
Message<String> message = new GenericMessage<String>("foo");
|
||||
this.input.send(message);
|
||||
Message<?> receive = this.output.receive(10000);
|
||||
assertNotNull(receive);
|
||||
assertEquals(1, this.adviceCalled.get());
|
||||
assertEquals(1, TestUtils.getPropertyValue(this.store, "metadata", Map.class).size());
|
||||
assertNotNull(this.store.get("foo"));
|
||||
String foo = this.store.get("foo");
|
||||
assertEquals("FOO", foo);
|
||||
|
||||
try {
|
||||
this.input.send(message);
|
||||
@@ -108,6 +121,18 @@ public class IdempotentReceiverIntegrationTests {
|
||||
assertEquals(1, TestUtils.getPropertyValue(store, "metadata", Map.class).size());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testIdempotentReceiverOnMethod() {
|
||||
TestUtils.getPropertyValue(this.store, "metadata", Map.class).clear();
|
||||
Message<String> message = new GenericMessage<String>("foo");
|
||||
this.annotatedMethodChannel.send(message);
|
||||
this.annotatedMethodChannel.send(message);
|
||||
|
||||
assertEquals(2, this.fooService.messages.size());
|
||||
assertTrue(this.fooService.messages.get(1).getHeaders().get(IntegrationMessageHeaderAccessor.DUPLICATE_MESSAGE,
|
||||
Boolean.class));
|
||||
}
|
||||
|
||||
@Configuration
|
||||
@EnableIntegration
|
||||
@EnableIntegrationMBeanExport(server = "mBeanServer")
|
||||
@@ -125,17 +150,21 @@ public class IdempotentReceiverIntegrationTests {
|
||||
|
||||
@Bean
|
||||
public IdempotentReceiverInterceptor idempotentReceiverInterceptor() {
|
||||
IdempotentReceiverInterceptor idempotentReceiverInterceptor =
|
||||
new IdempotentReceiverInterceptor(new MetadataStoreSelector(new MetadataKeyStrategy() {
|
||||
return new IdempotentReceiverInterceptor(new MetadataStoreSelector(new MessageProcessor<String>() {
|
||||
|
||||
@Override
|
||||
public String getKey(Message<?> message) {
|
||||
public String processMessage(Message<?> message) {
|
||||
return message.getPayload().toString();
|
||||
}
|
||||
|
||||
}, new MessageProcessor<String>() {
|
||||
|
||||
@Override
|
||||
public String processMessage(Message<?> message) {
|
||||
return message.getPayload().toString().toUpperCase();
|
||||
}
|
||||
|
||||
}, store()));
|
||||
idempotentReceiverInterceptor.setThrowExceptionOnRejection(true);
|
||||
return idempotentReceiverInterceptor;
|
||||
}
|
||||
|
||||
@Bean
|
||||
@@ -182,6 +211,29 @@ public class IdempotentReceiverIntegrationTests {
|
||||
};
|
||||
}
|
||||
|
||||
@Bean
|
||||
public MessageChannel annotatedMethodChannel() {
|
||||
return new DirectChannel();
|
||||
}
|
||||
|
||||
@Bean
|
||||
public FooService fooService() {
|
||||
return new FooService();
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@Component
|
||||
private static class FooService {
|
||||
|
||||
private List<Message<?>> messages = new ArrayList<Message<?>>();
|
||||
|
||||
@ServiceActivator(inputChannel = "annotatedMethodChannel")
|
||||
@IdempotentReceiver("idempotentReceiverInterceptor")
|
||||
public void handle(Message<?> message) {
|
||||
this.messages.add(message);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user