From a2cd79ea71b70b3b85c123bce109f66bd71b10da Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Fri, 5 Feb 2016 11:08:14 -0500 Subject: [PATCH] INT-3950: Fix `Aggregator` documentation JIRA: https://jira.spring.io/browse/INT-3950 Previously there was a mention of the `MessageGroupStore.expireMessageGroup(groupId)` which just doesn't existing in the Framework and never has been there. * Fix the documentation for the existing `MessageGroupStore.expireMessageGroups(timeout)`. Although the mention there of `Control Bus` requires to have `@ManagedOperation` on the method. * Add `@ManagedOperation` for the `MessageGroupStore.expireMessageGroups(timeout)` and confirm with the test-case: `AggregatorWithMessageStoreParserTests` * Fix the same docs in the XSD for `` * Fix other typos in the `aggregator.adoc` and `resequencer.adoc` Conflicts: spring-integration-core/src/main/java/org/springframework/integration/store/AbstractMessageGroupStore.java --- .../store/AbstractMessageGroupStore.java | 8 +++++--- .../integration/store/MessageGroupStore.java | 4 +++- ...atorWithMessageStoreParserTests-context.xml | 4 +++- .../AggregatorWithMessageStoreParserTests.java | 18 +++++++++++------- .../store/ConfigurableMongoDbMessageStore.java | 2 ++ src/reference/asciidoc/aggregator.adoc | 9 ++++----- src/reference/asciidoc/resequencer.adoc | 3 +-- 7 files changed, 29 insertions(+), 19 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractMessageGroupStore.java b/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractMessageGroupStore.java index 01e62ab088..414fc2bcff 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractMessageGroupStore.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractMessageGroupStore.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2015 the original author or authors. + * Copyright 2002-2016 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. You may obtain a copy of the License at @@ -26,6 +26,7 @@ import org.springframework.integration.support.DefaultMessageBuilderFactory; import org.springframework.integration.support.MessageBuilderFactory; import org.springframework.integration.support.utils.IntegrationUtils; import org.springframework.jmx.export.annotation.ManagedAttribute; +import org.springframework.jmx.export.annotation.ManagedOperation; import org.springframework.jmx.export.annotation.ManagedResource; import org.springframework.messaging.Message; @@ -103,11 +104,12 @@ public abstract class AbstractMessageGroupStore extends AbstractBatchingMessageG @Override public void registerMessageGroupExpiryCallback(MessageGroupCallback callback) { - expiryCallbacks.add(callback); + this.expiryCallbacks.add(callback); } @Override - public int expireMessageGroups(long timeout) { + @ManagedOperation + public synchronized int expireMessageGroups(long timeout) { int count = 0; long threshold = System.currentTimeMillis() - timeout; for (MessageGroup group : this) { diff --git a/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroupStore.java b/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroupStore.java index 66e79d23e3..2b92228108 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroupStore.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroupStore.java @@ -16,6 +16,7 @@ import java.util.Collection; import java.util.Iterator; import org.springframework.jmx.export.annotation.ManagedAttribute; +import org.springframework.jmx.export.annotation.ManagedOperation; import org.springframework.messaging.Message; /** @@ -98,6 +99,7 @@ public interface MessageGroupStore extends BasicMessageGroupStore { * * @see #registerMessageGroupExpiryCallback(MessageGroupCallback) */ + @ManagedOperation int expireMessageGroups(long timeout); /** @@ -142,7 +144,7 @@ public interface MessageGroupStore extends BasicMessageGroupStore { /** * Invoked when a MessageGroupStore expires a group. */ - public interface MessageGroupCallback { + interface MessageGroupCallback { void execute(MessageGroupStore messageGroupStore, MessageGroup group); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/AggregatorWithMessageStoreParserTests-context.xml b/spring-integration-core/src/test/java/org/springframework/integration/config/AggregatorWithMessageStoreParserTests-context.xml index 8a2781a4e0..c80a20a8ad 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/AggregatorWithMessageStoreParserTests-context.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/AggregatorWithMessageStoreParserTests-context.xml @@ -14,10 +14,12 @@ - + + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/AggregatorWithMessageStoreParserTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/AggregatorWithMessageStoreParserTests.java index 6eb6fdb24a..0fd26b82f4 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/AggregatorWithMessageStoreParserTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/AggregatorWithMessageStoreParserTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2010 the original author or authors. + * Copyright 2002-2016 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. @@ -26,31 +26,36 @@ import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.integration.store.MessageGroupStore; import org.springframework.integration.support.MessageBuilder; +import org.springframework.messaging.support.GenericMessage; import org.springframework.test.annotation.DirtiesContext; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; /** * @author Dave Syer + * @author Artem Bilan */ @ContextConfiguration @RunWith(SpringJUnit4ClassRunner.class) +@DirtiesContext public class AggregatorWithMessageStoreParserTests { - + @Autowired @Qualifier("input") private MessageChannel input; - + @Autowired private TestAggregatorBean aggregatorBean; - + @Autowired private MessageGroupStore messageGroupStore; + @Autowired + private MessageChannel controlBusChannel; + @Test @DirtiesContext public void testAggregation() { - input.send(createMessage("123", "id1", 3, 1, null)); assertEquals(1, messageGroupStore.getMessageGroup("id1").size()); input.send(createMessage("789", "id1", 3, 3, null)); @@ -67,12 +72,11 @@ public class AggregatorWithMessageStoreParserTests { @Test @DirtiesContext public void testExpiry() { - input.send(createMessage("123", "id1", 3, 1, null)); assertEquals(1, messageGroupStore.getMessageGroup("id1").size()); input.send(createMessage("456", "id1", 3, 2, null)); assertEquals(2, messageGroupStore.getMessageGroup("id1").size()); - messageGroupStore.expireMessageGroups(-10000); + this.controlBusChannel.send(new GenericMessage("@messageStore.expireMessageGroups(-10000)")); assertEquals("One and only one message should have been aggregated", 1, aggregatorBean .getAggregatedMessages().size()); Message aggregatedMessage = aggregatorBean.getAggregatedMessages().get("id1"); diff --git a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageStore.java b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageStore.java index f40e1a6983..f8be1e7ff5 100644 --- a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageStore.java +++ b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageStore.java @@ -37,6 +37,7 @@ import org.springframework.integration.store.MessageGroupStore; import org.springframework.integration.store.MessageStore; import org.springframework.integration.store.SimpleMessageGroup; import org.springframework.jmx.export.annotation.ManagedAttribute; +import org.springframework.jmx.export.annotation.ManagedOperation; import org.springframework.messaging.Message; import org.springframework.util.Assert; @@ -289,6 +290,7 @@ public class ConfigurableMongoDbMessageStore extends AbstractConfigurableMongoDb } @Override + @ManagedOperation public int expireMessageGroups(long timeout) { int count = 0; long threshold = System.currentTimeMillis() - timeout; diff --git a/src/reference/asciidoc/aggregator.adoc b/src/reference/asciidoc/aggregator.adoc index 30c2c00fc4..759948633c 100644 --- a/src/reference/asciidoc/aggregator.adoc +++ b/src/reference/asciidoc/aggregator.adoc @@ -342,13 +342,13 @@ _Optional_, by default a volatile in-memory store. <8> Indicates that expired messages should be aggregated and sent to the 'output-channel' or 'replyChannel' once their containing `MessageGroup` is expired (see `MessageGroupStore.expireMessageGroups(long)`). One way of expiring `MessageGroup` s is by configuring a `MessageGroupStoreReaper`. -However `MessageGroup` s can alternatively be expired by simply calling `MessageGroupStore.expireMessageGroup(groupId)`. +However `MessageGroup` s can alternatively be expired by simply calling `MessageGroupStore.expireMessageGroups(timeout)`. That could be accomplished via a Control Bus operation or by simply invoking that method if you have a reference to the `MessageGroupStore` instance. Otherwise by itself this attribute has no behavior. It only serves as an indicator of what to do (discard or send to the output/reply channel) with Messages that are still in the `MessageGroup` that is about to be expired. _Optional_. -_Default - 'false'_. -*NOTE:* This attribute is more properly 'send-partial-result-on-timeout' because the group may not actually expire if +_Default - false_. +*NOTE:* This attribute is more properly `send-partial-result-on-timeout` because the group may not actually expire if `expire-groups-upon-timeout` is set to `false`. @@ -367,8 +367,7 @@ _Optional_. <10> A reference to a bean that implements the message correlation (grouping) algorithm. The bean can be an implementation of the `CorrelationStrategy` interface or a POJO. In the latter case the correlation-strategy-method attribute must be defined as well. -_Optional (by default, the aggregator will use - the `IntegrationMessageHeaderAccessor.CORRELATION_ID` header) _. +_Optional (by default, the aggregator will use the `IntegrationMessageHeaderAccessor.CORRELATION_ID` header)_. diff --git a/src/reference/asciidoc/resequencer.adoc b/src/reference/asciidoc/resequencer.adoc index c870502280..04478b1a68 100644 --- a/src/reference/asciidoc/resequencer.adoc +++ b/src/reference/asciidoc/resequencer.adoc @@ -96,8 +96,7 @@ _Optional_. <9> A reference to a bean that implements the message correlation (grouping) algorithm. The bean can be an implementation of the `CorrelationStrategy` interface or a POJO. In the latter case the correlation-strategy-method attribute must be defined as well. -_Optional (by default, the aggregator will use - the `IntegrationMessageHeaderAccessor.CORRELATION_ID` header) _. +_Optional (by default, the aggregator will use the `IntegrationMessageHeaderAccessor.CORRELATION_ID` header)_.