diff --git a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java index ac7324918a..198058db7d 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2022 the original author or authors. + * Copyright 2002-2023 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. @@ -566,7 +566,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageP this.logger.trace(() -> "Adding message to group [ " + messageGroupToLog + "]"); messageGroup = store(correlationKey, message); - setGroupConditionIfAny(message, messageGroup); + messageGroup = setGroupConditionIfAny(message, messageGroup); if (this.releaseStrategy.canRelease(messageGroup)) { Collection> completedMessages = null; @@ -605,12 +605,19 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageP } } - private void setGroupConditionIfAny(Message message, MessageGroup messageGroup) { + private MessageGroup setGroupConditionIfAny(Message message, MessageGroup messageGroup) { + MessageGroup messageGroupToUse = messageGroup; + if (this.groupConditionSupplier != null) { - String condition = this.groupConditionSupplier.apply(message, messageGroup.getCondition()); - this.messageStore.setGroupCondition(messageGroup.getGroupId(), condition); - messageGroup.setCondition(condition); + String condition = this.groupConditionSupplier.apply(message, messageGroupToUse.getCondition()); + this.messageStore.setGroupCondition(messageGroupToUse.getGroupId(), condition); + messageGroupToUse = this.messageStore.getMessageGroup(messageGroupToUse.getGroupId()); + if (this.sequenceAware) { + messageGroupToUse = new SequenceAwareMessageGroup(messageGroupToUse); + } } + + return messageGroupToUse; } protected boolean isExpireGroupsUponCompletion() { diff --git a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageGroupStoreTests.java b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageGroupStoreTests.java index 62e2453e5f..d4c74e4707 100644 --- a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageGroupStoreTests.java +++ b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageGroupStoreTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2022 the original author or authors. + * Copyright 2002-2023 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. @@ -65,6 +65,29 @@ class ConfigurableMongoDbMessageGroupStoreTests extends AbstractMongoDbMessageGr super.testWithAggregatorWithShutdown("mongo-aggregator-configurable-config.xml"); } + @Test + void groupIsForceReleaseAfterTimeoutWhenGroupConditionIsSet() { + try (var context = new ClassPathXmlApplicationContext("mongo-aggregator-configurable-config.xml", getClass())) { + MessageChannel input = context.getBean("inputChannel", MessageChannel.class); + QueueChannel output = context.getBean("outputChannel", QueueChannel.class); + + Message message = MessageBuilder.withPayload("test") + .setSequenceNumber(1) + .setSequenceSize(10) + .setCorrelationId("test") + .build(); + + input.send(message); + + Message receive = output.receive(10_000); + + assertThat(receive) + .extracting("payload") + .asList() + .hasSize(1); + } + } + @Test @Disabled("The performance test. Enough slow. Also needs the release strategy changed to size() == 1000") void messageGroupStoreLazyLoadPerformance() { diff --git a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/mongo-aggregator-configurable-config.xml b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/mongo-aggregator-configurable-config.xml index 552b5ec778..4a407d8429 100644 --- a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/mongo-aggregator-configurable-config.xml +++ b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/mongo-aggregator-configurable-config.xml @@ -11,6 +11,8 @@