GH-9558: Expose BarrierSpec.discardChannel & triggerTimeout

Fixes: #9558
Issue link: https://github.com/spring-projects/spring-integration/issues/9558

**Auto-cherry-pick to `6.3.x` & `6.2.x`**
This commit is contained in:
Artem Bilan
2024-10-18 16:06:06 -04:00
parent 15e914b75b
commit e1cebafa5c
4 changed files with 72 additions and 5 deletions

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2015-2020 the original author or authors.
* Copyright 2015-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.
@@ -27,6 +27,7 @@ import org.springframework.integration.handler.AbstractReplyProducingMessageHand
import org.springframework.integration.handler.DiscardingMessageHandler;
import org.springframework.integration.handler.MessageTriggerAction;
import org.springframework.integration.store.SimpleMessageGroup;
import org.springframework.lang.Nullable;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.MessageHandlingException;
@@ -170,7 +171,8 @@ public class BarrierMessageHandler extends AbstractReplyProducingMessageHandler
}
/**
* Set the name of the channel to which late arriving trigger messages are sent.
* Set the name of the channel to which late arriving trigger messages are sent,
* or request message does not arrive in time.
* @param discardChannelName the discard channel.
* @since 5.0
*/
@@ -179,7 +181,8 @@ public class BarrierMessageHandler extends AbstractReplyProducingMessageHandler
}
/**
* Set the channel to which late arriving trigger messages are sent.
* Set the channel to which late arriving trigger messages are sent,
* or request message does not arrive in time.
* @param discardChannel the discard channel.
* @since 5.0
*/
@@ -188,8 +191,11 @@ public class BarrierMessageHandler extends AbstractReplyProducingMessageHandler
}
/**
* Return the discard message channel for trigger action message.
* @return a discard message channel.
* @since 5.0
*/
@Nullable
@Override
public MessageChannel getDiscardChannel() {
String channelName = this.discardChannelName;

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2016-2023 the original author or authors.
* Copyright 2016-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.
@@ -25,6 +25,8 @@ import org.springframework.integration.aggregator.DefaultAggregatingMessageGroup
import org.springframework.integration.aggregator.HeaderAttributeCorrelationStrategy;
import org.springframework.integration.aggregator.MessageGroupProcessor;
import org.springframework.integration.config.ConsumerEndpointFactoryBean;
import org.springframework.lang.Nullable;
import org.springframework.messaging.MessageChannel;
import org.springframework.util.Assert;
/**
@@ -43,6 +45,15 @@ public class BarrierSpec extends ConsumerEndpointSpec<BarrierSpec, BarrierMessag
private CorrelationStrategy correlationStrategy =
new HeaderAttributeCorrelationStrategy(IntegrationMessageHeaderAccessor.CORRELATION_ID);
@Nullable
private MessageChannel discardChannel;
@Nullable
private String discardChannelName;
@Nullable
private Long triggerTimeout;
protected BarrierSpec(long timeout) {
super(null);
this.timeout = timeout;
@@ -60,9 +71,57 @@ public class BarrierSpec extends ConsumerEndpointSpec<BarrierSpec, BarrierMessag
return this;
}
/**
* Set the channel to which late arriving trigger messages are sent,
* or request message does not arrive in time.
* @param discardChannel the message channel for discarded triggers.
* @return the spec
* @since 6.2.10
*/
public BarrierSpec discardChannel(@Nullable MessageChannel discardChannel) {
this.discardChannel = discardChannel;
return this;
}
/**
* Set the channel bean name to which late arriving trigger messages are sent,
* or request message does not arrive in time.
* @param discardChannelName the message channel for discarded triggers.
* @return the spec
* @since 6.2.10
*/
public BarrierSpec discardChannel(@Nullable String discardChannelName) {
this.discardChannelName = discardChannelName;
return this;
}
/**
* Set the timeout in milliseconds when waiting for a request message.
* @param triggerTimeout the timeout in milliseconds when waiting for a request message.
* @return the spec
* @since 6.2.10
*/
public BarrierSpec triggerTimeout(long triggerTimeout) {
this.triggerTimeout = triggerTimeout;
return this;
}
@Override
public Tuple2<ConsumerEndpointFactoryBean, BarrierMessageHandler> doGet() {
this.handler = new BarrierMessageHandler(this.timeout, this.outputProcessor, this.correlationStrategy);
if (this.triggerTimeout == null) {
this.handler = new BarrierMessageHandler(this.timeout, this.outputProcessor, this.correlationStrategy);
}
else {
this.handler =
new BarrierMessageHandler(this.timeout, this.triggerTimeout, this.outputProcessor,
this.correlationStrategy);
}
if (this.discardChannel != null) {
this.handler.setDiscardChannel(this.discardChannel);
}
else if (this.discardChannelName != null) {
this.handler.setDiscardChannelName(this.discardChannelName);
}
return super.doGet();
}

View File

@@ -300,6 +300,7 @@ public class CorrelationHandlerTests {
return f -> f
.barrier(10000, b -> b
.correlationStrategy(new HeaderAttributeCorrelationStrategy(BARRIER))
.discardChannel("nullChannel")
.outputProcessor(g ->
g.getMessages()
.stream()