Fix compatibility with the latest SF

* Mostly changes are related to the `TaskScheduler` and `Trigger` APIs
* Migrate to `micrometer-tracing` dependency
* Rework `SocketTestUtils` to use a `InetAddress.getLocalHost()`
for more stability and performance on Windows
* Fix docs for new `PeriodicTrigger` API
This commit is contained in:
Artem Bilan
2022-07-08 17:36:15 -04:00
parent 64aa4d5348
commit e0f137905a
46 changed files with 565 additions and 586 deletions

View File

@@ -84,8 +84,8 @@ ext {
lettuceVersion = '6.1.8.RELEASE'
log4jVersion = '2.17.2'
mailVersion = '2.0.1'
micrometerVersion = '1.10.0-M2'
micrometerTracingVersion = '1.0.0-M5'
micrometerVersion = '1.10.0-SNAPSHOT'
micrometerTracingVersion = '1.0.0-SNAPSHOT'
mockitoVersion = '4.5.1'
mongoDriverVersion = '4.6.0'
mysqlVersion = '8.0.28'
@@ -502,7 +502,7 @@ project('spring-integration-core') {
optionalApi "com.jayway.jsonpath:json-path:$jsonpathVersion"
optionalApi "com.esotericsoftware:kryo:$kryoVersion"
optionalApi 'io.micrometer:micrometer-core'
optionalApi ('io.micrometer:micrometer-tracing-api') {
optionalApi ('io.micrometer:micrometer-tracing') {
exclude group: 'aopalliance'
}
optionalApi "io.github.resilience4j:resilience4j-ratelimiter:$resilience4jVersion"

View File

@@ -17,6 +17,7 @@
package org.springframework.integration.aggregator;
import java.time.Duration;
import java.time.Instant;
import java.util.Collection;
import java.util.Collections;
import java.util.Comparator;
@@ -650,7 +651,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageP
removeEmptyGroupAfterTimeout(groupId, timeout);
}
}, new Date(System.currentTimeMillis() + timeout));
}, Instant.now().plusMillis(timeout));
this.logger.debug(() -> "Schedule empty MessageGroup [ " + groupId + "] for removal.");
this.expireGroupScheduledFutures.put(groupId, scheduledFuture);
@@ -687,7 +688,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageP
"] is rescheduled by the reason of: ");
scheduleGroupToForceComplete(groupId);
}
}, startTime);
}, startTime.toInstant());
this.logger.debug(() -> "Schedule MessageGroup [ " + messageGroup + "] to 'forceComplete'.");
this.expireGroupScheduledFutures.put(UUIDConverter.getUUID(groupId), scheduledFuture);

View File

@@ -16,6 +16,7 @@
package org.springframework.integration.channel;
import java.time.Instant;
import java.util.Date;
import java.util.Iterator;
import java.util.Map;
@@ -70,14 +71,14 @@ public class DefaultHeaderChannelRegistry extends IntegrationObjectSupport
private volatile boolean explicitlyStopped;
/**
* Constructs a registry with the default delay for channel expiry.
* Construct a registry with the default delay for channel expiry.
*/
public DefaultHeaderChannelRegistry() {
this(DEFAULT_REAPER_DELAY);
}
/**
* Constructs a registry with the provided delay (milliseconds) for
* Construct a registry with the provided delay (milliseconds) for
* channel expiry.
* @param reaperDelay the delay in milliseconds.
*/
@@ -125,7 +126,7 @@ public class DefaultHeaderChannelRegistry extends IntegrationObjectSupport
Assert.notNull(getTaskScheduler(), "a task scheduler is required");
this.reaperScheduledFuture =
getTaskScheduler()
.schedule(this, new Date(System.currentTimeMillis() + this.reaperDelay));
.schedule(this, Instant.now().plusMillis(this.reaperDelay));
this.running = true;
}

View File

@@ -18,6 +18,7 @@ package org.springframework.integration.config.annotation;
import java.lang.annotation.Annotation;
import java.lang.reflect.Method;
import java.time.Duration;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
@@ -516,10 +517,10 @@ public abstract class AbstractMethodAnnotationPostProcessor<T extends Annotation
else if (StringUtils.hasText(fixedDelayValue)) {
Assert.state(!StringUtils.hasText(fixedRateValue),
"The '@Poller' 'fixedDelay' attribute is mutually exclusive with other attributes.");
trigger = new PeriodicTrigger(Long.parseLong(fixedDelayValue));
trigger = new PeriodicTrigger(Duration.ofMillis(Long.parseLong(fixedDelayValue)));
}
else if (StringUtils.hasText(fixedRateValue)) {
trigger = new PeriodicTrigger(Long.parseLong(fixedRateValue));
trigger = new PeriodicTrigger(Duration.ofMillis(Long.parseLong(fixedRateValue)));
((PeriodicTrigger) trigger).setFixedRate(true);
}
//'Trigger' can be null. 'PollingConsumer' does fallback to the 'new PeriodicTrigger(10)'.

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2016-2019 the original author or authors.
* Copyright 2016-2022 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.
@@ -16,6 +16,8 @@
package org.springframework.integration.dsl;
import java.time.Duration;
import java.time.temporal.ChronoUnit;
import java.util.TimeZone;
import java.util.concurrent.TimeUnit;
@@ -53,7 +55,7 @@ public final class PollerFactory {
}
public PollerSpec fixedRate(long period, TimeUnit timeUnit) {
return Pollers.fixedRate(period, timeUnit);
return Pollers.fixedRate(Duration.of(period, timeUnit.toChronoUnit()));
}
public PollerSpec fixedRate(long period, long initialDelay) {
@@ -61,15 +63,17 @@ public final class PollerFactory {
}
public PollerSpec fixedDelay(long period, TimeUnit timeUnit, long initialDelay) {
return Pollers.fixedDelay(period, timeUnit, initialDelay);
ChronoUnit chronoUnit = timeUnit.toChronoUnit();
return Pollers.fixedDelay(Duration.of(period, chronoUnit), Duration.of(initialDelay, chronoUnit));
}
public PollerSpec fixedRate(long period, TimeUnit timeUnit, long initialDelay) {
return Pollers.fixedRate(period, timeUnit, initialDelay);
ChronoUnit chronoUnit = timeUnit.toChronoUnit();
return Pollers.fixedRate(Duration.of(period, chronoUnit), Duration.of(initialDelay, chronoUnit));
}
public PollerSpec fixedDelay(long period, TimeUnit timeUnit) {
return Pollers.fixedDelay(period, timeUnit);
return Pollers.fixedDelay(Duration.of(period, timeUnit.toChronoUnit()));
}
public PollerSpec fixedDelay(long period, long initialDelay) {

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2016-2019 the original author or authors.
* Copyright 2016-2022 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.
@@ -17,6 +17,7 @@
package org.springframework.integration.dsl;
import java.time.Duration;
import java.time.temporal.ChronoUnit;
import java.util.TimeZone;
import java.util.concurrent.TimeUnit;
@@ -40,56 +41,74 @@ public final class Pollers {
return new PollerSpec(trigger);
}
public static PollerSpec fixedRate(Duration period) {
return fixedRate(period.toMillis());
public static PollerSpec fixedRate(long period) {
return fixedRate(Duration.ofMillis(period));
}
public static PollerSpec fixedRate(long period) {
public static PollerSpec fixedRate(Duration period) {
return fixedRate(period, null);
}
/**
* @deprecated since 6.0 in favor of {@link #fixedRate(Duration)}
*/
@Deprecated(since = "6.0", forRemoval = true)
public static PollerSpec fixedRate(long period, TimeUnit timeUnit) {
return fixedRate(period, timeUnit, 0);
return fixedRate(Duration.of(period, timeUnit.toChronoUnit()));
}
public static PollerSpec fixedRate(Duration period, Duration initialDelay) {
return fixedRate(period.toMillis(), initialDelay.toMillis());
return periodicTrigger(period, true, initialDelay);
}
public static PollerSpec fixedRate(long period, long initialDelay) {
return periodicTrigger(period, null, true, initialDelay);
return fixedRate(Duration.ofMillis(period), Duration.ofMillis(initialDelay));
}
/**
* @deprecated since 6.0 in favor of {@link #fixedRate(Duration, Duration)}
*/
@Deprecated(since = "6.0", forRemoval = true)
public static PollerSpec fixedRate(long period, TimeUnit timeUnit, long initialDelay) {
return periodicTrigger(period, timeUnit, true, initialDelay);
ChronoUnit chronoUnit = timeUnit.toChronoUnit();
return fixedRate(Duration.of(period, chronoUnit), Duration.of(initialDelay, chronoUnit));
}
public static PollerSpec fixedDelay(Duration period) {
return fixedDelay(period.toMillis());
}
public static PollerSpec fixedDelay(long period) {
return fixedDelay(period, null);
}
public static PollerSpec fixedDelay(Duration period, Duration initialDelay) {
return fixedDelay(period.toMillis(), initialDelay.toMillis());
public static PollerSpec fixedDelay(long period) {
return fixedDelay(Duration.ofMillis(period));
}
public static PollerSpec fixedDelay(Duration period, Duration initialDelay) {
return periodicTrigger(period, false, initialDelay);
}
/**
* @deprecated since 6.0 in favor of {@link #fixedDelay(Duration)}
*/
@Deprecated(since = "6.0", forRemoval = true)
public static PollerSpec fixedDelay(long period, TimeUnit timeUnit) {
return fixedDelay(period, timeUnit, 0);
return fixedDelay(Duration.of(period, timeUnit.toChronoUnit()));
}
public static PollerSpec fixedDelay(long period, long initialDelay) {
return periodicTrigger(period, null, false, initialDelay);
return fixedDelay(Duration.ofMillis(period), Duration.ofMillis(initialDelay));
}
/**
* @deprecated since 6.0 in favor of {@link #fixedDelay(Duration, Duration)}
*/
@Deprecated(since = "6.0", forRemoval = true)
public static PollerSpec fixedDelay(long period, TimeUnit timeUnit, long initialDelay) {
return periodicTrigger(period, timeUnit, false, initialDelay);
ChronoUnit chronoUnit = timeUnit.toChronoUnit();
return fixedDelay(Duration.of(period, chronoUnit), Duration.of(initialDelay, chronoUnit));
}
private static PollerSpec periodicTrigger(long period, TimeUnit timeUnit, boolean fixedRate, long initialDelay) {
PeriodicTrigger periodicTrigger = new PeriodicTrigger(period, timeUnit);
private static PollerSpec periodicTrigger(Duration period, boolean fixedRate, Duration initialDelay) {
PeriodicTrigger periodicTrigger = new PeriodicTrigger(period);
periodicTrigger.setFixedRate(fixedRate);
periodicTrigger.setInitialDelay(initialDelay);
return new PollerSpec(periodicTrigger);

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2021 the original author or authors.
* Copyright 2002-2022 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.
@@ -17,8 +17,8 @@
package org.springframework.integration.endpoint;
import java.time.Duration;
import java.time.Instant;
import java.util.Collection;
import java.util.Date;
import java.util.HashSet;
import java.util.List;
import java.util.concurrent.Callable;
@@ -97,7 +97,7 @@ public abstract class AbstractPollingEndpoint extends AbstractEndpoint implement
private ClassLoader beanClassLoader = ClassUtils.getDefaultClassLoader();
private Trigger trigger = new PeriodicTrigger(DEFAULT_POLLING_PERIOD);
private Trigger trigger = new PeriodicTrigger(Duration.ofMillis(DEFAULT_POLLING_PERIOD));
private ErrorHandler errorHandler;
@@ -139,7 +139,7 @@ public abstract class AbstractPollingEndpoint extends AbstractEndpoint implement
}
public void setTrigger(Trigger trigger) {
this.trigger = (trigger != null ? trigger : new PeriodicTrigger(DEFAULT_POLLING_PERIOD));
this.trigger = (trigger != null ? trigger : new PeriodicTrigger(Duration.ofMillis(DEFAULT_POLLING_PERIOD)));
}
public void setAdviceChain(List<Advice> adviceChain) {
@@ -148,7 +148,7 @@ public abstract class AbstractPollingEndpoint extends AbstractEndpoint implement
/**
* Configure a cap for messages to poll from the source per scheduling cycle.
* A negative number means retrieve unlimited messages until the {@code MessageSource} returns {@code null}.
* A negative number means retrieve unlimited messages until the {@code MessageSource} returns {@code null}.
* Zero means do not poll for any records - it
* can be considered as pausing if 'maxMessagesPerPoll' is later changed to a non-zero value.
* The polling cycle may exit earlier if the source returns null for the current receive call.
@@ -174,6 +174,7 @@ public abstract class AbstractPollingEndpoint extends AbstractEndpoint implement
public void setTransactionSynchronizationFactory(TransactionSynchronizationFactory
transactionSynchronizationFactory) {
this.transactionSynchronizationFactory = transactionSynchronizationFactory;
}
@@ -200,7 +201,7 @@ public abstract class AbstractPollingEndpoint extends AbstractEndpoint implement
* Return true if this advice should be applied only to the {@link #receiveMessage()} operation
* rather than the whole poll.
* @param advice The advice.
* @return true to only advise the receive operation.
* @return true to only advise the {@code receive} operation.
*/
protected boolean isReceiveOnlyAdvice(Advice advice) {
return advice instanceof ReceiveMessageAdvice;
@@ -226,7 +227,7 @@ public abstract class AbstractPollingEndpoint extends AbstractEndpoint implement
}
this.appliedAdvices.clear();
this.appliedAdvices.addAll(chain);
if (!(isSyncExecutor())) {
if (!isSyncExecutor()) {
logger.warn(() -> getComponentName() + ": A task executor is supplied and " + chain.size()
+ "ReceiveMessageAdvice(s) is/are provided. If an advice mutates the source, such "
+ "mutations are not thread safe and could cause unexpected results, especially with "
@@ -358,11 +359,10 @@ public abstract class AbstractPollingEndpoint extends AbstractEndpoint implement
return Flux
.<Duration>generate(sink -> {
Date date = this.trigger.nextExecutionTime(triggerContext);
Instant date = this.trigger.nextExecution(triggerContext);
if (date != null) {
triggerContext.update(date, null, null);
long millis = date.getTime() - System.currentTimeMillis();
sink.next(Duration.ofMillis(millis));
sink.next(Duration.ofMillis(date.toEpochMilli() - System.currentTimeMillis()));
}
else {
sink.complete();
@@ -371,8 +371,8 @@ public abstract class AbstractPollingEndpoint extends AbstractEndpoint implement
.concatMap(duration ->
Mono.delay(duration)
.doOnNext(l ->
triggerContext.update(triggerContext.lastScheduledExecutionTime(),
new Date(), null))
triggerContext.update(
triggerContext.lastScheduledExecution(), Instant.now(), null))
.flatMapMany(l ->
Flux
.defer(() -> {
@@ -400,9 +400,9 @@ public abstract class AbstractPollingEndpoint extends AbstractEndpoint implement
.subscribeOn(Schedulers.fromExecutor(this.taskExecutor))
.doOnComplete(() ->
triggerContext
.update(triggerContext.lastScheduledExecutionTime(),
triggerContext.lastActualExecutionTime(),
new Date())
.update(triggerContext.lastScheduledExecution(),
triggerContext.lastActualExecution(),
Instant.now())
)), 0)
.repeat(this::isActive)
.doOnSubscribe(subs -> this.subscription = subs);
@@ -412,9 +412,9 @@ public abstract class AbstractPollingEndpoint extends AbstractEndpoint implement
try {
return this.pollingTask.call();
}
catch (Exception e) {
if (e instanceof MessagingException) { // NOSONAR
throw (MessagingException) e;
catch (Exception ex) {
if (ex instanceof MessagingException) { // NOSONAR
throw (MessagingException) ex;
}
else {
Message<?> failedMessage = null;
@@ -424,7 +424,7 @@ public abstract class AbstractPollingEndpoint extends AbstractEndpoint implement
failedMessage = ((IntegrationResourceHolder) resource).getMessage();
}
}
throw new MessagingException(failedMessage, e); // NOSONAR (null failedMessage)
throw new MessagingException(failedMessage, ex); // NOSONAR (null failedMessage)
}
}
finally {
@@ -509,7 +509,7 @@ public abstract class AbstractPollingEndpoint extends AbstractEndpoint implement
protected abstract void handleMessage(Message<?> message);
/**
* Return a resource (MessageSource etc) to bind when using transaction
* Return a resource (MessageSource etc.) to bind when using transaction
* synchronization.
* @return The resource, or null if transaction synchronization is not required.
*/
@@ -537,9 +537,7 @@ public abstract class AbstractPollingEndpoint extends AbstractEndpoint implement
if (synchronization != null) {
TransactionSynchronizationManager.registerSynchronization(synchronization);
if (synchronization instanceof IntegrationResourceHolderSynchronization) {
IntegrationResourceHolderSynchronization integrationSynchronization =
((IntegrationResourceHolderSynchronization) synchronization);
if (synchronization instanceof IntegrationResourceHolderSynchronization integrationSynchronization) {
integrationSynchronization.setShouldUnbindAtCompletion(false);
if (!TransactionSynchronizationManager.hasResource(resource)) {
@@ -550,8 +548,7 @@ public abstract class AbstractPollingEndpoint extends AbstractEndpoint implement
}
Object resourceHolder = TransactionSynchronizationManager.getResource(resource);
if (resourceHolder instanceof IntegrationResourceHolder) {
IntegrationResourceHolder integrationResourceHolder = (IntegrationResourceHolder) resourceHolder;
if (resourceHolder instanceof IntegrationResourceHolder integrationResourceHolder) {
if (key != null) {
integrationResourceHolder.addAttribute(key, resource);
}

View File

@@ -17,6 +17,7 @@
package org.springframework.integration.handler;
import java.io.Serializable;
import java.time.Instant;
import java.util.Collection;
import java.util.Collections;
import java.util.Date;
@@ -426,7 +427,7 @@ public class DelayHandler extends AbstractReplyProducingMessageHandler implement
}
Runnable releaseTask = releaseTaskForMessage(delayedMessage);
Date startTime = new Date(messageWrapper.getRequestDate() + delay);
Instant startTime = Instant.ofEpochMilli(messageWrapper.getRequestDate() + delay);
if (TransactionSynchronizationManager.isSynchronizationActive() &&
TransactionSynchronizationManager.isActualTransactionActive()) {
@@ -481,9 +482,9 @@ public class DelayHandler extends AbstractReplyProducingMessageHandler implement
this.releaseHandler.handleMessage(message);
this.deliveries.remove(identity);
}
catch (Exception e) {
catch (Exception ex) {
if (getErrorChannel() != null) {
ErrorMessage errorMessage = new ErrorMessage(e,
ErrorMessage errorMessage = new ErrorMessage(ex,
Collections.singletonMap(IntegrationMessageHeaderAccessor.DELIVERY_ATTEMPT,
new AtomicInteger(this.deliveries.get(identity).get() + 1)),
message);
@@ -502,9 +503,9 @@ public class DelayHandler extends AbstractReplyProducingMessageHandler implement
}
}
else {
logger.debug(e, () -> "Release flow threw an exception for message: " + message);
logger.debug(ex, () -> "Release flow threw an exception for message: " + message);
if (!rescheduleForRetry(message, identity)) {
throw e; // there might be an error handler on the scheduler
throw ex; // there might be an error handler on the scheduler
}
}
}
@@ -600,7 +601,7 @@ public class DelayHandler extends AbstractReplyProducingMessageHandler implement
/**
* Handle {@link ContextRefreshedEvent} to invoke
* {@link #reschedulePersistedMessages} as late as possible after application context
* startup. Also it checks {@link #initialized} to ignore other
* startup. Also, it checks {@link #initialized} to ignore other
* {@link ContextRefreshedEvent}s which may be published in the 'parent-child'
* contexts, e.g. in the Spring-MVC applications.
* @param event - {@link ContextRefreshedEvent} which occurs after Application context

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2015-2019 the original author or authors.
* Copyright 2015-2022 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.
@@ -16,7 +16,7 @@
package org.springframework.integration.util;
import java.util.Date;
import java.time.Instant;
import org.springframework.scheduling.Trigger;
import org.springframework.scheduling.TriggerContext;
@@ -29,6 +29,8 @@ import org.springframework.util.Assert;
* invoked.
*
* @author Gary Russell
* @author Artem Bilan
*
* @since 4.3
*
*/
@@ -65,12 +67,12 @@ public class CompoundTrigger implements Trigger {
}
@Override
public Date nextExecutionTime(TriggerContext triggerContext) {
public Instant nextExecution(TriggerContext triggerContext) {
if (this.override != null) {
return this.override.nextExecutionTime(triggerContext);
return this.override.nextExecution(triggerContext);
}
else {
return this.primary.nextExecutionTime(triggerContext);
return this.primary.nextExecution(triggerContext);
}
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2020 the original author or authors.
* Copyright 2002-2022 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.
@@ -17,7 +17,8 @@
package org.springframework.integration.util;
import java.time.Duration;
import java.util.Date;
import java.time.Instant;
import java.util.Objects;
import org.springframework.scheduling.Trigger;
import org.springframework.scheduling.TriggerContext;
@@ -91,6 +92,7 @@ public class DynamicPeriodicTrigger implements Trigger {
* @since 5.1
*/
public void setDuration(Duration duration) {
Assert.notNull(duration, "duration must not be null");
this.duration = duration;
}
@@ -124,30 +126,19 @@ public class DynamicPeriodicTrigger implements Trigger {
/**
* Return the time after which a task should run again.
* @param triggerContext the trigger context to determine the previous state of schedule.
* @return the the next schedule date.
* @return the next schedule instant.
*/
@Override
public Date nextExecutionTime(TriggerContext triggerContext) {
Date lastScheduled = triggerContext.lastScheduledExecutionTime();
public Instant nextExecution(TriggerContext triggerContext) {
Instant lastScheduled = triggerContext.lastScheduledExecution();
if (lastScheduled == null) {
return new Date(System.currentTimeMillis() + this.initialDuration.toMillis());
return Instant.now().plus(this.initialDuration);
}
else if (this.fixedRate) {
return new Date(lastScheduled.getTime() + this.duration.toMillis());
return lastScheduled.plus(this.duration);
}
return
new Date(triggerContext.lastCompletionTime().getTime() + // NOSONAR never null here
this.duration.toMillis());
}
@Override
public int hashCode() {
final int prime = 31;
int result = 1;
result = prime * result + ((this.duration == null) ? 0 : this.duration.hashCode());
result = prime * result + (this.fixedRate ? 1231 : 1237); // NOSONAR
result = prime * result + ((this.initialDuration == null) ? 0 : this.initialDuration.hashCode());
return result;
triggerContext.lastCompletion().plus(this.duration); // NOSONAR never null here
}
@Override
@@ -155,30 +146,18 @@ public class DynamicPeriodicTrigger implements Trigger {
if (this == obj) {
return true;
}
if (obj == null) {
if (obj == null || getClass() != obj.getClass()) {
return false;
}
if (getClass() != obj.getClass()) {
return false;
}
DynamicPeriodicTrigger other = (DynamicPeriodicTrigger) obj;
if (this.duration == null) {
if (other.duration != null) {
return false;
}
}
else if (!this.duration.equals(other.duration)) {
return false;
}
if (this.fixedRate != other.fixedRate) {
return false;
}
if (this.initialDuration == null) {
return other.initialDuration == null;
}
else {
return this.initialDuration.equals(other.initialDuration);
}
DynamicPeriodicTrigger that = (DynamicPeriodicTrigger) obj;
return this.fixedRate == that.fixedRate
&& Objects.equals(this.initialDuration, that.initialDuration)
&& Objects.equals(this.duration, that.duration);
}
@Override
public int hashCode() {
return Objects.hash(this.initialDuration, this.duration, this.fixedRate);
}
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2019 the original author or authors.
* Copyright 2002-2022 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.
@@ -19,10 +19,13 @@ package org.springframework.integration.bus;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.Mockito.mock;
import java.time.Duration;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import org.junit.Test;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.BeanFactory;
import org.springframework.context.support.ClassPathXmlApplicationContext;
@@ -51,13 +54,24 @@ import org.springframework.scheduling.support.PeriodicTrigger;
*/
public class ApplicationContextMessageBusTests {
private TestApplicationContext context;
@BeforeEach
void setup() {
this.context = TestUtils.createTestApplicationContext();
}
@AfterEach
void tearDown() {
this.context.close();
}
@Test
public void endpointRegistrationWithInputChannelReference() {
TestApplicationContext context = TestUtils.createTestApplicationContext();
QueueChannel sourceChannel = new QueueChannel();
QueueChannel targetChannel = new QueueChannel();
context.registerChannel("sourceChannel", sourceChannel);
context.registerChannel("targetChannel", targetChannel);
this.context.registerChannel("sourceChannel", sourceChannel);
this.context.registerChannel("targetChannel", targetChannel);
Message<String> message = MessageBuilder.withPayload("test")
.setReplyChannelName("targetChannel").build();
sourceChannel.send(message);
@@ -68,29 +82,26 @@ public class ApplicationContextMessageBusTests {
return message;
}
};
handler.setBeanFactory(context);
handler.setBeanFactory(this.context);
handler.afterPropertiesSet();
PollingConsumer endpoint = new PollingConsumer(sourceChannel, handler);
endpoint.setBeanFactory(mock(BeanFactory.class));
context.registerEndpoint("testEndpoint", endpoint);
context.refresh();
this.context.registerEndpoint("testEndpoint", endpoint);
this.context.refresh();
Message<?> result = targetChannel.receive(10000);
assertThat(result.getPayload()).isEqualTo("test");
context.close();
}
@Test
public void channelsWithoutHandlers() {
TestApplicationContext context = TestUtils.createTestApplicationContext();
QueueChannel sourceChannel = new QueueChannel();
context.registerChannel("sourceChannel", sourceChannel);
this.context.registerChannel("sourceChannel", sourceChannel);
sourceChannel.send(new GenericMessage<>("test"));
QueueChannel targetChannel = new QueueChannel();
context.registerChannel("targetChannel", targetChannel);
context.refresh();
this.context.registerChannel("targetChannel", targetChannel);
this.context.refresh();
Message<?> result = targetChannel.receive(10);
assertThat(result).isNull();
context.close();
}
@Test
@@ -108,7 +119,6 @@ public class ApplicationContextMessageBusTests {
@Test
public void exactlyOneConsumerReceivesPointToPointMessage() {
TestApplicationContext context = TestUtils.createTestApplicationContext();
QueueChannel inputChannel = new QueueChannel();
QueueChannel outputChannel1 = new QueueChannel();
QueueChannel outputChannel2 = new QueueChannel();
@@ -126,28 +136,26 @@ public class ApplicationContextMessageBusTests {
return message;
}
};
context.registerChannel("input", inputChannel);
context.registerChannel("output1", outputChannel1);
context.registerChannel("output2", outputChannel2);
this.context.registerChannel("input", inputChannel);
this.context.registerChannel("output1", outputChannel1);
this.context.registerChannel("output2", outputChannel2);
handler1.setOutputChannel(outputChannel1);
handler2.setOutputChannel(outputChannel2);
PollingConsumer endpoint1 = new PollingConsumer(inputChannel, handler1);
endpoint1.setBeanFactory(mock(BeanFactory.class));
PollingConsumer endpoint2 = new PollingConsumer(inputChannel, handler2);
endpoint2.setBeanFactory(mock(BeanFactory.class));
context.registerEndpoint("testEndpoint1", endpoint1);
context.registerEndpoint("testEndpoint2", endpoint2);
context.refresh();
this.context.registerEndpoint("testEndpoint1", endpoint1);
this.context.registerEndpoint("testEndpoint2", endpoint2);
this.context.refresh();
inputChannel.send(new GenericMessage<>("testing"));
Message<?> message1 = outputChannel1.receive(10000);
Message<?> message2 = outputChannel2.receive(0);
context.close();
assertThat(message1 == null ^ message2 == null).as("exactly one message should be null").isTrue();
}
@Test
public void bothConsumersReceivePublishSubscribeMessage() throws InterruptedException {
TestApplicationContext context = TestUtils.createTestApplicationContext();
PublishSubscribeChannel inputChannel = new PublishSubscribeChannel();
QueueChannel outputChannel1 = new QueueChannel();
QueueChannel outputChannel2 = new QueueChannel();
@@ -168,43 +176,40 @@ public class ApplicationContextMessageBusTests {
return message;
}
};
context.registerChannel("input", inputChannel);
context.registerChannel("output1", outputChannel1);
context.registerChannel("output2", outputChannel2);
this.context.registerChannel("input", inputChannel);
this.context.registerChannel("output1", outputChannel1);
this.context.registerChannel("output2", outputChannel2);
handler1.setOutputChannel(outputChannel1);
handler2.setOutputChannel(outputChannel2);
EventDrivenConsumer endpoint1 = new EventDrivenConsumer(inputChannel, handler1);
EventDrivenConsumer endpoint2 = new EventDrivenConsumer(inputChannel, handler2);
context.registerEndpoint("testEndpoint1", endpoint1);
context.registerEndpoint("testEndpoint2", endpoint2);
context.refresh();
inputChannel.send(new GenericMessage<String>("testing"));
this.context.registerEndpoint("testEndpoint1", endpoint1);
this.context.registerEndpoint("testEndpoint2", endpoint2);
this.context.refresh();
inputChannel.send(new GenericMessage<>("testing"));
latch.await(500, TimeUnit.MILLISECONDS);
assertThat(latch.getCount()).as("both handlers should have been invoked").isEqualTo(0);
Message<?> message1 = outputChannel1.receive(500);
Message<?> message2 = outputChannel2.receive(500);
context.close();
assertThat(message1).as("both handlers should have replied to the message").isNotNull();
assertThat(message2).as("both handlers should have replied to the message").isNotNull();
}
@Test
public void errorChannelWithFailedDispatch() throws InterruptedException {
TestApplicationContext context = TestUtils.createTestApplicationContext();
QueueChannel errorChannel = new QueueChannel();
QueueChannel outputChannel = new QueueChannel();
context.registerChannel("errorChannel", errorChannel);
this.context.registerChannel("errorChannel", errorChannel);
CountDownLatch latch = new CountDownLatch(1);
SourcePollingChannelAdapter channelAdapter = new SourcePollingChannelAdapter();
channelAdapter.setSource(new FailingSource(latch));
PollerMetadata pollerMetadata = new PollerMetadata();
pollerMetadata.setTrigger(new PeriodicTrigger(1000));
pollerMetadata.setTrigger(new PeriodicTrigger(Duration.ofSeconds(1)));
channelAdapter.setOutputChannel(outputChannel);
context.registerEndpoint("testChannel", channelAdapter);
context.refresh();
this.context.registerEndpoint("testChannel", channelAdapter);
this.context.refresh();
latch.await(2000, TimeUnit.MILLISECONDS);
Message<?> message = errorChannel.receive(5000);
context.close();
assertThat(outputChannel.receive(100)).isNull();
assertThat(message).as("message should not be null").isNotNull();
assertThat(message instanceof ErrorMessage).isTrue();
@@ -214,9 +219,8 @@ public class ApplicationContextMessageBusTests {
@Test
public void consumerSubscribedToErrorChannel() throws InterruptedException {
TestApplicationContext context = TestUtils.createTestApplicationContext();
QueueChannel errorChannel = new QueueChannel();
context.registerChannel(IntegrationContextUtils.ERROR_CHANNEL_BEAN_NAME, errorChannel);
this.context.registerChannel(IntegrationContextUtils.ERROR_CHANNEL_BEAN_NAME, errorChannel);
final CountDownLatch latch = new CountDownLatch(1);
AbstractReplyProducingMessageHandler handler = new AbstractReplyProducingMessageHandler() {
@@ -228,22 +232,15 @@ public class ApplicationContextMessageBusTests {
};
PollingConsumer endpoint = new PollingConsumer(errorChannel, handler);
endpoint.setBeanFactory(mock(BeanFactory.class));
context.registerEndpoint("testEndpoint", endpoint);
context.refresh();
this.context.registerEndpoint("testEndpoint", endpoint);
this.context.refresh();
errorChannel.send(new ErrorMessage(new RuntimeException("test-exception")));
latch.await(1000, TimeUnit.MILLISECONDS);
assertThat(latch.getCount()).as("handler should have received error message").isEqualTo(0);
context.close();
}
private static class FailingSource implements MessageSource<Object> {
private final CountDownLatch latch;
FailingSource(CountDownLatch latch) {
this.latch = latch;
}
private record FailingSource(CountDownLatch latch) implements MessageSource<Object> {
@Override
public Message<Object> receive() {

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2021 the original author or authors.
* Copyright 2002-2022 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 static org.mockito.Mockito.verify;
import static org.mockito.Mockito.verifyNoInteractions;
import static org.mockito.Mockito.when;
import java.time.Duration;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.CountDownLatch;
@@ -85,7 +86,7 @@ public class SourcePollingChannelAdapterFactoryBeanTests {
adviceApplied.set(true);
return invocation.proceed();
});
pollerMetadata.setTrigger(new PeriodicTrigger(5000));
pollerMetadata.setTrigger(new PeriodicTrigger(Duration.ofSeconds(5)));
pollerMetadata.setMaxMessagesPerPoll(1);
pollerMetadata.setAdviceChain(adviceChain);
factoryBean.setPollerMetadata(pollerMetadata);
@@ -115,7 +116,7 @@ public class SourcePollingChannelAdapterFactoryBeanTests {
adviceApplied.set(true);
return invocation.proceed();
});
pollerMetadata.setTrigger(new PeriodicTrigger(5000));
pollerMetadata.setTrigger(new PeriodicTrigger(Duration.ofSeconds(5)));
pollerMetadata.setMaxMessagesPerPoll(1);
final AtomicInteger count = new AtomicInteger();
final MethodInterceptor txAdvice = mock(MethodInterceptor.class);
@@ -232,7 +233,7 @@ public class SourcePollingChannelAdapterFactoryBeanTests {
SourcePollingChannelAdapter pollingChannelAdapter = new SourcePollingChannelAdapter();
pollingChannelAdapter.setTaskScheduler(taskScheduler);
pollingChannelAdapter.setSource(() -> new GenericMessage<>("test"));
pollingChannelAdapter.setTrigger(new PeriodicTrigger(1));
pollingChannelAdapter.setTrigger(new PeriodicTrigger(Duration.ofMillis(1)));
pollingChannelAdapter.setMaxMessagesPerPoll(0);
QueueChannel outputChannel = new QueueChannel();
pollingChannelAdapter.setOutputChannel(outputChannel);

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2019 the original author or authors.
* Copyright 2002-2022 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.
@@ -16,7 +16,8 @@
package org.springframework.integration.config;
import java.util.Date;
import java.time.Duration;
import java.time.Instant;
import java.util.concurrent.ScheduledFuture;
import org.springframework.scheduling.TaskScheduler;
@@ -24,30 +25,38 @@ import org.springframework.scheduling.Trigger;
/**
* @author Mark Fisher
* @author Artem Bilan
*/
public class StubTaskScheduler implements TaskScheduler {
@Override
public ScheduledFuture<?> schedule(Runnable task, Trigger trigger) {
return null;
}
public ScheduledFuture<?> schedule(Runnable task, Date startTime) {
@Override
public ScheduledFuture<?> schedule(Runnable task, Instant startTime) {
return null;
}
public ScheduledFuture<?> scheduleAtFixedRate(Runnable task, long period) {
@Override
public ScheduledFuture<?> scheduleAtFixedRate(Runnable task, Instant startTime, Duration period) {
return null;
}
public ScheduledFuture<?> scheduleAtFixedRate(Runnable task, Date startTime, long period) {
@Override
public ScheduledFuture<?> scheduleAtFixedRate(Runnable task, Duration period) {
return null;
}
public ScheduledFuture<?> scheduleWithFixedDelay(Runnable task, long delay) {
@Override
public ScheduledFuture<?> scheduleWithFixedDelay(Runnable task, Instant startTime, Duration delay) {
return null;
}
public ScheduledFuture<?> scheduleWithFixedDelay(Runnable task, Date startTime, long delay) {
@Override
public ScheduledFuture<?> scheduleWithFixedDelay(Runnable task, Duration delay) {
return null;
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2016-2019 the original author or authors.
* Copyright 2016-2022 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.
@@ -16,17 +16,19 @@
package org.springframework.integration.config;
import java.util.Date;
import java.time.Instant;
import org.springframework.scheduling.Trigger;
import org.springframework.scheduling.TriggerContext;
/**
* @author Marius Bogoevici
* @author Artem Bilan
*/
public class TestTrigger implements Trigger {
public Date nextExecutionTime(TriggerContext triggerContext) {
@Override
public Instant nextExecution(TriggerContext triggerContext) {
throw new UnsupportedOperationException();
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2021 the original author or authors.
* Copyright 2002-2022 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.
@@ -18,6 +18,7 @@ package org.springframework.integration.config.xml;
import static org.assertj.core.api.Assertions.assertThat;
import java.time.Duration;
import java.util.Map;
import org.junit.jupiter.api.Test;
@@ -57,7 +58,7 @@ public class InboundChannelAdapterExpressionTests {
Trigger trigger = TestUtils.getPropertyValue(adapter, "trigger", Trigger.class);
assertThat(trigger.getClass()).isEqualTo(PeriodicTrigger.class);
DirectFieldAccessor triggerAccessor = new DirectFieldAccessor(trigger);
assertThat(triggerAccessor.getPropertyValue("period")).isEqualTo(1234L);
assertThat(triggerAccessor.getPropertyValue("period")).isEqualTo(Duration.ofMillis(1234));
assertThat(triggerAccessor.getPropertyValue("fixedRate")).isEqualTo(Boolean.FALSE);
assertThat(adapterAccessor.getPropertyValue("outputChannel"))
.isEqualTo(this.context.getBean("fixedDelayChannel"));
@@ -74,7 +75,7 @@ public class InboundChannelAdapterExpressionTests {
Trigger trigger = TestUtils.getPropertyValue(adapter, "trigger", Trigger.class);
assertThat(trigger.getClass()).isEqualTo(PeriodicTrigger.class);
DirectFieldAccessor triggerAccessor = new DirectFieldAccessor(trigger);
assertThat(triggerAccessor.getPropertyValue("period")).isEqualTo(5678L);
assertThat(triggerAccessor.getPropertyValue("period")).isEqualTo(Duration.ofMillis(5678));
assertThat(triggerAccessor.getPropertyValue("fixedRate")).isEqualTo(Boolean.TRUE);
assertThat(adapterAccessor.getPropertyValue("outputChannel"))
.isEqualTo(this.context.getBean("fixedRateChannel"));

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2019 the original author or authors.
* Copyright 2002-2022 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.
@@ -18,23 +18,22 @@ package org.springframework.integration.config.xml;
import static org.assertj.core.api.Assertions.assertThat;
import org.junit.Test;
import org.junit.runner.RunWith;
import java.time.Duration;
import org.junit.jupiter.api.Test;
import org.springframework.beans.DirectFieldAccessor;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.integration.endpoint.SourcePollingChannelAdapter;
import org.springframework.integration.test.util.TestUtils;
import org.springframework.scheduling.Trigger;
import org.springframework.scheduling.support.PeriodicTrigger;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
import org.springframework.test.context.junit.jupiter.SpringJUnitConfig;
/**
* @author Mark Fisher
* @author Artem Bilan
*/
@ContextConfiguration
@RunWith(SpringJUnit4ClassRunner.class)
@SpringJUnitConfig
public class InboundChannelAdapterWithDefaultPollerTests {
@Autowired
@@ -44,10 +43,10 @@ public class InboundChannelAdapterWithDefaultPollerTests {
@Test
public void verifyDefaultPollerInUse() {
Trigger trigger = TestUtils.getPropertyValue(adapter, "trigger", Trigger.class);
assertThat(trigger.getClass()).isEqualTo(PeriodicTrigger.class);
DirectFieldAccessor triggerAccessor = new DirectFieldAccessor(trigger);
assertThat(triggerAccessor.getPropertyValue("period")).isEqualTo(12345L);
assertThat(triggerAccessor.getPropertyValue("fixedRate")).isEqualTo(Boolean.TRUE);
assertThat(trigger).isInstanceOf(PeriodicTrigger.class);
PeriodicTrigger periodicTrigger = (PeriodicTrigger) trigger;
assertThat(periodicTrigger.getPeriodDuration()).isEqualTo(Duration.ofMillis(12345));
assertThat(periodicTrigger.isFixedRate()).isTrue();
}
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2019 the original author or authors.
* Copyright 2002-2022 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.
@@ -18,53 +18,45 @@ package org.springframework.integration.config.xml;
import static org.assertj.core.api.Assertions.assertThat;
import org.junit.Test;
import org.junit.runner.RunWith;
import java.time.Duration;
import org.junit.jupiter.api.Test;
import org.springframework.beans.DirectFieldAccessor;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.ApplicationContext;
import org.springframework.integration.scheduling.PollerMetadata;
import org.springframework.scheduling.Trigger;
import org.springframework.scheduling.support.PeriodicTrigger;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
import org.springframework.test.context.junit.jupiter.SpringJUnitConfig;
/**
* @author Marius Bogoevici
* @author Artem Bilan
*/
@ContextConfiguration
@RunWith(SpringJUnit4ClassRunner.class)
@SpringJUnitConfig
public class IntervalTriggerParserTests {
@Autowired
ApplicationContext context;
PollerMetadata pollerWithFixedRateAttribute;
@Autowired
PollerMetadata pollerWithFixedDelayAttribute;
@Test
public void testFixedRateTrigger() {
Object poller = context.getBean("pollerWithFixedRateAttribute");
assertThat(poller.getClass()).isEqualTo(PollerMetadata.class);
PollerMetadata metadata = (PollerMetadata) poller;
Trigger trigger = metadata.getTrigger();
assertThat(trigger.getClass()).isEqualTo(PeriodicTrigger.class);
DirectFieldAccessor accessor = new DirectFieldAccessor(trigger);
Boolean fixedRate = (Boolean) accessor.getPropertyValue("fixedRate");
Long period = (Long) accessor.getPropertyValue("period");
assertThat(true).isEqualTo(fixedRate);
assertThat(period.longValue()).isEqualTo(36L);
Trigger trigger = this.pollerWithFixedRateAttribute.getTrigger();
assertThat(trigger).isInstanceOf(PeriodicTrigger.class);
PeriodicTrigger periodicTrigger = (PeriodicTrigger) trigger;
assertThat(periodicTrigger.getPeriodDuration()).isEqualTo(Duration.ofMillis(36));
assertThat(periodicTrigger.isFixedRate()).isTrue();
}
@Test
public void testFixedDelayTrigger() {
Object poller = context.getBean("pollerWithFixedDelayAttribute");
assertThat(poller.getClass()).isEqualTo(PollerMetadata.class);
PollerMetadata metadata = (PollerMetadata) poller;
Trigger trigger = metadata.getTrigger();
assertThat(trigger.getClass()).isEqualTo(PeriodicTrigger.class);
DirectFieldAccessor accessor = new DirectFieldAccessor(trigger);
Boolean fixedRate = (Boolean) accessor.getPropertyValue("fixedRate");
Long period = (Long) accessor.getPropertyValue("period");
assertThat(false).isEqualTo(fixedRate);
assertThat(period.longValue()).isEqualTo(37L);
Trigger trigger = this.pollerWithFixedDelayAttribute.getTrigger();
assertThat(trigger).isInstanceOf(PeriodicTrigger.class);
PeriodicTrigger periodicTrigger = (PeriodicTrigger) trigger;
assertThat(periodicTrigger.getPeriodDuration()).isEqualTo(Duration.ofMillis(37));
assertThat(periodicTrigger.isFixedRate()).isFalse();
}
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2019 the original author or authors.
* Copyright 2002-2022 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.
@@ -17,12 +17,13 @@
package org.springframework.integration.config.xml;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatExceptionOfType;
import java.time.temporal.ChronoUnit;
import java.util.HashMap;
import java.util.concurrent.TimeUnit;
import org.aopalliance.aop.Advice;
import org.junit.Test;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.parsing.BeanDefinitionParsingException;
import org.springframework.context.support.ClassPathXmlApplicationContext;
@@ -62,16 +63,18 @@ public class PollerParserTests {
context.close();
}
@Test(expected = BeanDefinitionParsingException.class)
@Test
public void multipleDefaultPollers() {
new ClassPathXmlApplicationContext(
"multipleDefaultPollers.xml", PollerParserTests.class).close();
assertThatExceptionOfType(BeanDefinitionParsingException.class)
.isThrownBy(() ->
new ClassPathXmlApplicationContext("multipleDefaultPollers.xml", PollerParserTests.class));
}
@Test(expected = BeanDefinitionParsingException.class)
@Test
public void topLevelPollerWithoutId() {
new ClassPathXmlApplicationContext(
"topLevelPollerWithoutId.xml", PollerParserTests.class).close();
assertThatExceptionOfType(BeanDefinitionParsingException.class)
.isThrownBy(() ->
new ClassPathXmlApplicationContext("topLevelPollerWithoutId.xml", PollerParserTests.class));
}
@Test
@@ -89,8 +92,9 @@ public class PollerParserTests {
assertThat(metadata.getAdviceChain().get(2)).isSameAs(context.getBean("adviceBean3"));
Advice txAdvice = metadata.getAdviceChain().get(3);
assertThat(txAdvice.getClass()).isEqualTo(TransactionInterceptor.class);
TransactionAttributeSource transactionAttributeSource = ((TransactionInterceptor) txAdvice).getTransactionAttributeSource();
assertThat(transactionAttributeSource.getClass()).isEqualTo(NameMatchTransactionAttributeSource.class);
TransactionAttributeSource transactionAttributeSource =
((TransactionInterceptor) txAdvice).getTransactionAttributeSource();
assertThat(transactionAttributeSource).isInstanceOf(NameMatchTransactionAttributeSource.class);
@SuppressWarnings("rawtypes")
HashMap nameMap = TestUtils.getPropertyValue(transactionAttributeSource, "nameMap", HashMap.class);
assertThat(nameMap.size()).isEqualTo(1);
@@ -108,7 +112,7 @@ public class PollerParserTests {
PollerMetadata metadata = (PollerMetadata) poller;
assertThat(metadata.getReceiveTimeout()).isEqualTo(1234);
PeriodicTrigger trigger = (PeriodicTrigger) metadata.getTrigger();
assertThat(TestUtils.getPropertyValue(trigger, "timeUnit").toString()).isEqualTo(TimeUnit.SECONDS.toString());
assertThat(TestUtils.getPropertyValue(trigger, "chronoUnit")).isEqualTo(ChronoUnit.SECONDS);
context.close();
}
@@ -123,22 +127,25 @@ public class PollerParserTests {
context.close();
}
@Test(expected = BeanDefinitionParsingException.class)
@Test
public void pollerWithCronTriggerAndTimeUnit() {
new ClassPathXmlApplicationContext(
"cronTriggerWithTimeUnit-fail.xml", PollerParserTests.class).close();
assertThatExceptionOfType(BeanDefinitionParsingException.class)
.isThrownBy(() ->
new ClassPathXmlApplicationContext("cronTriggerWithTimeUnit-fail.xml", PollerParserTests.class));
}
@Test(expected = BeanDefinitionParsingException.class)
@Test
public void topLevelPollerWithRef() {
new ClassPathXmlApplicationContext(
"defaultPollerWithRef.xml", PollerParserTests.class).close();
assertThatExceptionOfType(BeanDefinitionParsingException.class)
.isThrownBy(() ->
new ClassPathXmlApplicationContext("defaultPollerWithRef.xml", PollerParserTests.class));
}
@Test(expected = BeanDefinitionParsingException.class)
@Test
public void pollerWithCronAndFixedDelay() {
new ClassPathXmlApplicationContext(
"pollerWithCronAndFixedDelay.xml", PollerParserTests.class).close();
assertThatExceptionOfType(BeanDefinitionParsingException.class)
.isThrownBy(() ->
new ClassPathXmlApplicationContext("pollerWithCronAndFixedDelay.xml", PollerParserTests.class));
}
}

View File

@@ -28,6 +28,7 @@ import java.lang.annotation.ElementType;
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
import java.lang.annotation.Target;
import java.time.Duration;
import java.util.Arrays;
import java.util.Date;
import java.util.List;
@@ -62,6 +63,7 @@ import org.springframework.context.annotation.ImportResource;
import org.springframework.context.annotation.Lazy;
import org.springframework.context.expression.EnvironmentAccessor;
import org.springframework.context.expression.MapAccessor;
import org.springframework.core.annotation.AliasFor;
import org.springframework.core.convert.converter.Converter;
import org.springframework.core.log.LogAccessor;
import org.springframework.core.serializer.support.SerializingConverter;
@@ -319,8 +321,9 @@ public class EnableIntegrationTests {
Trigger trigger = TestUtils.getPropertyValue(this.serviceActivatorEndpoint, "trigger", Trigger.class);
assertThat(trigger).isInstanceOf(PeriodicTrigger.class);
assertThat(TestUtils.getPropertyValue(trigger, "period")).isEqualTo(100L);
assertThat(TestUtils.getPropertyValue(trigger, "fixedRate", Boolean.class)).isFalse();
PeriodicTrigger periodicTrigger = (PeriodicTrigger) trigger;
assertThat(periodicTrigger.getPeriodDuration()).isEqualTo(Duration.ofMillis(100));
assertThat(periodicTrigger.isFixedRate()).isFalse();
assertThat(this.annotationTestService.isRunning()).isTrue();
LogAccessor logger = spy(TestUtils.getPropertyValue(this.serviceActivatorEndpoint, "logger", LogAccessor.class));
@@ -343,8 +346,9 @@ public class EnableIntegrationTests {
trigger = TestUtils.getPropertyValue(this.serviceActivatorEndpoint1, "trigger", Trigger.class);
assertThat(trigger).isInstanceOf(PeriodicTrigger.class);
assertThat(TestUtils.getPropertyValue(trigger, "period")).isEqualTo(100L);
assertThat(TestUtils.getPropertyValue(trigger, "fixedRate", Boolean.class)).isTrue();
periodicTrigger = (PeriodicTrigger) trigger;
assertThat(periodicTrigger.getPeriodDuration()).isEqualTo(Duration.ofMillis(100));
assertThat(periodicTrigger.isFixedRate()).isTrue();
trigger = TestUtils.getPropertyValue(this.serviceActivatorEndpoint2, "trigger", Trigger.class);
assertThat(trigger).isInstanceOf(CronTrigger.class);
@@ -352,19 +356,22 @@ public class EnableIntegrationTests {
trigger = TestUtils.getPropertyValue(this.serviceActivatorEndpoint3, "trigger", Trigger.class);
assertThat(trigger).isInstanceOf(PeriodicTrigger.class);
assertThat(TestUtils.getPropertyValue(trigger, "period")).isEqualTo(11L);
assertThat(TestUtils.getPropertyValue(trigger, "fixedRate", Boolean.class)).isFalse();
periodicTrigger = (PeriodicTrigger) trigger;
assertThat(periodicTrigger.getPeriodDuration()).isEqualTo(Duration.ofMillis(11));
assertThat(periodicTrigger.isFixedRate()).isFalse();
trigger = TestUtils.getPropertyValue(this.serviceActivatorEndpoint4, "trigger", Trigger.class);
assertThat(trigger).isInstanceOf(PeriodicTrigger.class);
assertThat(TestUtils.getPropertyValue(trigger, "period")).isEqualTo(1000L);
assertThat(TestUtils.getPropertyValue(trigger, "fixedRate", Boolean.class)).isFalse();
periodicTrigger = (PeriodicTrigger) trigger;
assertThat(periodicTrigger.getPeriodDuration()).isEqualTo(Duration.ofSeconds(1));
assertThat(periodicTrigger.isFixedRate()).isFalse();
assertThat(trigger).isSameAs(this.myTrigger);
trigger = TestUtils.getPropertyValue(this.transformer, "trigger", Trigger.class);
assertThat(trigger).isInstanceOf(PeriodicTrigger.class);
assertThat(TestUtils.getPropertyValue(trigger, "period")).isEqualTo(10L);
assertThat(TestUtils.getPropertyValue(trigger, "fixedRate", Boolean.class)).isFalse();
periodicTrigger = (PeriodicTrigger) trigger;
assertThat(periodicTrigger.getPeriodDuration()).isEqualTo(Duration.ofMillis(10));
assertThat(periodicTrigger.isFixedRate()).isFalse();
this.input.send(MessageBuilder.withPayload("Foo").build());
@@ -588,7 +595,7 @@ public class EnableIntegrationTests {
assertThat(TestUtils.getPropertyValue(consumer, "handler.outputChannelName")).isEqualTo("annOutput");
assertThat(TestUtils.getPropertyValue(consumer,
"handler.adviceChain", List.class).get(0)).isSameAs(context.getBean("annAdvice"));
assertThat(TestUtils.getPropertyValue(consumer, "trigger.period")).isEqualTo(1000L);
assertThat(TestUtils.getPropertyValue(consumer, "trigger.period")).isEqualTo(Duration.ofSeconds(1));
consumer = this.context.getBean("annotationTestService.annCount1.serviceActivator",
PollingConsumer.class);
@@ -599,7 +606,7 @@ public class EnableIntegrationTests {
assertThat(TestUtils.getPropertyValue(consumer, "handler.outputChannel.beanName")).isEqualTo("annOutput");
assertThat(TestUtils.getPropertyValue(consumer,
"handler.adviceChain", List.class).get(0)).isSameAs(context.getBean("annAdvice1"));
assertThat(TestUtils.getPropertyValue(consumer, "trigger.period")).isEqualTo(2000L);
assertThat(TestUtils.getPropertyValue(consumer, "trigger.period")).isEqualTo(Duration.ofSeconds(2));
consumer = this.context.getBean("annotationTestService.annCount2.serviceActivator",
PollingConsumer.class);
@@ -609,7 +616,7 @@ public class EnableIntegrationTests {
assertThat(TestUtils.getPropertyValue(consumer, "handler.outputChannelName")).isEqualTo("annOutput");
assertThat(TestUtils.getPropertyValue(consumer,
"handler.adviceChain", List.class).get(0)).isSameAs(context.getBean("annAdvice"));
assertThat(TestUtils.getPropertyValue(consumer, "trigger.period")).isEqualTo(1000L);
assertThat(TestUtils.getPropertyValue(consumer, "trigger.period")).isEqualTo(Duration.ofSeconds(1));
// Tests when the channel is in a "middle" annotation
consumer = this.context.getBean("annotationTestService.annCount5.serviceActivator", PollingConsumer.class);
@@ -619,7 +626,7 @@ public class EnableIntegrationTests {
assertThat(TestUtils.getPropertyValue(consumer, "handler.outputChannelName")).isEqualTo("annOutput");
assertThat(TestUtils.getPropertyValue(consumer,
"handler.adviceChain", List.class).get(0)).isSameAs(context.getBean("annAdvice"));
assertThat(TestUtils.getPropertyValue(consumer, "trigger.period")).isEqualTo(1000L);
assertThat(TestUtils.getPropertyValue(consumer, "trigger.period")).isEqualTo(Duration.ofSeconds(1));
consumer = this.context.getBean("annotationTestService.annAgg1.aggregator", PollingConsumer.class);
assertThat(TestUtils.getPropertyValue(consumer, "autoStartup", Boolean.class)).isFalse();
@@ -627,7 +634,7 @@ public class EnableIntegrationTests {
assertThat(TestUtils.getPropertyValue(consumer, "inputChannel")).isSameAs(context.getBean("annInput"));
assertThat(TestUtils.getPropertyValue(consumer, "handler.outputChannelName")).isEqualTo("annOutput");
assertThat(TestUtils.getPropertyValue(consumer, "handler.discardChannelName")).isEqualTo("annOutput");
assertThat(TestUtils.getPropertyValue(consumer, "trigger.period")).isEqualTo(1000L);
assertThat(TestUtils.getPropertyValue(consumer, "trigger.period")).isEqualTo(Duration.ofSeconds(1));
assertThat(TestUtils.getPropertyValue(consumer, "handler.messagingTemplate.sendTimeout")).isEqualTo(-1L);
assertThat(TestUtils.getPropertyValue(consumer, "handler.sendPartialResultOnExpiry", Boolean.class)).isFalse();
@@ -637,7 +644,7 @@ public class EnableIntegrationTests {
assertThat(TestUtils.getPropertyValue(consumer, "inputChannel")).isSameAs(context.getBean("annInput"));
assertThat(TestUtils.getPropertyValue(consumer, "handler.outputChannelName")).isEqualTo("annOutput");
assertThat(TestUtils.getPropertyValue(consumer, "handler.discardChannelName")).isEqualTo("annOutput");
assertThat(TestUtils.getPropertyValue(consumer, "trigger.period")).isEqualTo(1000L);
assertThat(TestUtils.getPropertyValue(consumer, "trigger.period")).isEqualTo(Duration.ofSeconds(1));
assertThat(TestUtils.getPropertyValue(consumer, "handler.messagingTemplate.sendTimeout")).isEqualTo(75L);
assertThat(TestUtils.getPropertyValue(consumer, "handler.sendPartialResultOnExpiry", Boolean.class)).isTrue();
}
@@ -852,7 +859,7 @@ public class EnableIntegrationTests {
@Bean
public Trigger myTrigger() {
return new PeriodicTrigger(1000L);
return new PeriodicTrigger(Duration.ofSeconds(1));
}
@Bean
@@ -1193,14 +1200,14 @@ public class EnableIntegrationTests {
@Bean(name = PollerMetadata.DEFAULT_POLLER)
public PollerMetadata defaultPoller() {
PollerMetadata pollerMetadata = new PollerMetadata();
pollerMetadata.setTrigger(new PeriodicTrigger(10));
pollerMetadata.setTrigger(new PeriodicTrigger(Duration.ofMillis(10)));
return pollerMetadata;
}
@Bean
public PollerMetadata myPoller() {
PollerMetadata pollerMetadata = new PollerMetadata();
pollerMetadata.setTrigger(new PeriodicTrigger(11));
pollerMetadata.setTrigger(new PeriodicTrigger(Duration.ofMillis(11)));
return pollerMetadata;
}
@@ -1342,7 +1349,7 @@ public class EnableIntegrationTests {
assertThat(message.getHeaders().get("foo")).isEqualTo("FOO");
assertThat(message.getHeaders()).containsKey("calledMethod");
assertThat(message.getHeaders().get("calledMethod")).isEqualTo("echo");
return handle(message.getPayload()) + Arrays.asList(new Throwable().getStackTrace()).toString();
return handle(message.getPayload()) + Arrays.asList(new Throwable().getStackTrace());
}
@Transformer(inputChannel = "gatewayChannel2")
@@ -1352,7 +1359,7 @@ public class EnableIntegrationTests {
assertThat(message.getHeaders().get("foo")).isEqualTo("FOO");
assertThat(message.getHeaders()).containsKey("calledMethod");
assertThat(message.getHeaders().get("calledMethod")).isEqualTo("echo2");
return handle(message.getPayload()) + "2" + Arrays.asList(new Throwable().getStackTrace()).toString();
return handle(message.getPayload()) + "2" + Arrays.asList(new Throwable().getStackTrace());
}
@MyInboundChannelAdapter1
@@ -1490,6 +1497,7 @@ public class EnableIntegrationTests {
defaultHeaders = @GatewayHeader(name = "foo", value = "FOO"))
public @interface TestMessagingGateway {
@AliasFor(annotation = MessagingGateway.class, attribute = "defaultRequestChannel")
String defaultRequestChannel() default "";
}
@@ -1499,6 +1507,7 @@ public class EnableIntegrationTests {
@TestMessagingGateway(defaultRequestChannel = "gatewayChannel2")
public @interface TestMessagingGateway2 {
@AliasFor(annotation = TestMessagingGateway.class, attribute = "defaultRequestChannel")
String defaultRequestChannel() default "";
}
@@ -1513,16 +1522,22 @@ public class EnableIntegrationTests {
poller = @Poller(fixedDelay = "1000"))
public @interface MyServiceActivator {
@AliasFor(annotation = ServiceActivator.class, attribute = "inputChannel")
String inputChannel() default "";
@AliasFor(annotation = ServiceActivator.class, attribute = "outputChannel")
String outputChannel() default "";
@AliasFor(annotation = ServiceActivator.class, attribute = "adviceChain")
String[] adviceChain() default { };
@AliasFor(annotation = ServiceActivator.class, attribute = "autoStartup")
String autoStartup() default "";
@AliasFor(annotation = ServiceActivator.class, attribute = "phase")
String phase() default "";
@AliasFor(annotation = ServiceActivator.class, attribute = "poller")
Poller poller() default @Poller(ValueConstants.DEFAULT_NONE);
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2019 the original author or authors.
* Copyright 2019-2022 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.
@@ -19,7 +19,6 @@ package org.springframework.integration.dsl;
import static org.assertj.core.api.Assertions.assertThat;
import java.time.Duration;
import java.util.concurrent.TimeUnit;
import org.junit.jupiter.api.Test;
@@ -27,6 +26,8 @@ import org.springframework.scheduling.support.PeriodicTrigger;
/**
* @author Gary Russell
* @author Artem Bilan
*
* @since 5.1.4
*
*/
@@ -35,24 +36,20 @@ public class PollersTests {
@Test
public void testDurations() {
PeriodicTrigger trigger = (PeriodicTrigger) Pollers.fixedDelay(Duration.ofMinutes(1L)).get().getTrigger();
assertThat(trigger.getPeriod()).isEqualTo(60_000L);
assertThat(trigger.getTimeUnit()).isEqualTo(TimeUnit.MILLISECONDS);
assertThat(trigger.getPeriodDuration()).isEqualTo(Duration.ofSeconds(60));
assertThat(trigger.isFixedRate()).isFalse();
trigger = (PeriodicTrigger) Pollers.fixedDelay(Duration.ofMinutes(1L), Duration.ofSeconds(10L))
.get().getTrigger();
assertThat(trigger.getPeriod()).isEqualTo(60_000L);
assertThat(trigger.getInitialDelay()).isEqualTo(10_000L);
assertThat(trigger.getTimeUnit()).isEqualTo(TimeUnit.MILLISECONDS);
assertThat(trigger.getPeriodDuration()).isEqualTo(Duration.ofSeconds(60));
assertThat(trigger.getInitialDelayDuration()).isEqualTo(Duration.ofSeconds(10));
assertThat(trigger.isFixedRate()).isFalse();
trigger = (PeriodicTrigger) Pollers.fixedRate(Duration.ofMinutes(1L)).get().getTrigger();
assertThat(trigger.getPeriod()).isEqualTo(60_000L);
assertThat(trigger.getTimeUnit()).isEqualTo(TimeUnit.MILLISECONDS);
assertThat(trigger.getPeriodDuration()).isEqualTo(Duration.ofSeconds(60));
assertThat(trigger.isFixedRate()).isTrue();
trigger = (PeriodicTrigger) Pollers.fixedRate(Duration.ofMinutes(1L), Duration.ofSeconds(10L))
.get().getTrigger();
assertThat(trigger.getPeriod()).isEqualTo(60_000L);
assertThat(trigger.getInitialDelay()).isEqualTo(10_000L);
assertThat(trigger.getTimeUnit()).isEqualTo(TimeUnit.MILLISECONDS);
assertThat(trigger.getPeriodDuration()).isEqualTo(Duration.ofSeconds(60));
assertThat(trigger.getInitialDelayDuration()).isEqualTo(Duration.ofSeconds(10));
assertThat(trigger.isFixedRate()).isTrue();
}

View File

@@ -18,9 +18,9 @@ package org.springframework.integration.dsl.flowservices;
import static org.assertj.core.api.Assertions.assertThat;
import java.time.Instant;
import java.util.Collection;
import java.util.Collections;
import java.util.Date;
import java.util.List;
import java.util.Optional;
import java.util.concurrent.atomic.AtomicReference;
@@ -165,15 +165,15 @@ public class FlowServiceTests {
@Component
public static class MyFlowAdapter extends IntegrationFlowAdapter {
private final AtomicReference<Date> executionDate = new AtomicReference<>(new Date());
private final AtomicReference<Instant> executionDate = new AtomicReference<>(Instant.now());
private Date nextExecutionTime(TriggerContext triggerContext) {
private Instant nextExecution(TriggerContext triggerContext) {
return this.executionDate.getAndSet(null);
}
@Override
protected IntegrationFlowDefinition<?> buildFlow() {
return fromSupplier(this::messageSource, e -> e.poller(p -> p.trigger(this::nextExecutionTime)))
return fromSupplier(this::messageSource, e -> e.poller(p -> p.trigger(this::nextExecution)))
.split(this, null, e -> e.applySequence(false))
.transform(this)
.aggregate(a -> a.processor(this, null))

View File

@@ -22,6 +22,7 @@ import static org.assertj.core.api.Assertions.assertThatIllegalArgumentException
import static org.assertj.core.api.Assertions.assertThatIllegalStateException;
import java.lang.reflect.Method;
import java.time.Instant;
import java.util.Arrays;
import java.util.Date;
import java.util.List;
@@ -539,7 +540,7 @@ public class ManualFlowTests {
private static class MyFlowAdapter extends IntegrationFlowAdapter {
private final AtomicReference<Date> nextExecutionTime = new AtomicReference<>(new Date());
private final AtomicReference<Instant> nextExecutionTime = new AtomicReference<>(Instant.now());
@Override
protected IntegrationFlowDefinition<?> buildFlow() {

View File

@@ -19,9 +19,9 @@ package org.springframework.integration.dsl.reactivestreams;
import static org.assertj.core.api.Assertions.assertThat;
import java.time.Duration;
import java.time.Instant;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Date;
import java.util.List;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
@@ -109,11 +109,11 @@ public class ReactiveStreamsTests {
CountDownLatch latch = new CountDownLatch(6);
Disposable disposable =
Flux.from(this.publisher)
.map(m -> m.getPayload().toUpperCase())
.subscribe(p -> {
results.add(p);
latch.countDown();
});
.map(m -> m.getPayload().toUpperCase())
.subscribe(p -> {
results.add(p);
latch.countDown();
});
assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
String[] strings = results.toArray(new String[0]);
assertThat(strings).isEqualTo(new String[]{ "A", "B", "C", "D", "E", "F" });
@@ -251,7 +251,7 @@ public class ReactiveStreamsTests {
public Publisher<Message<String>> reactiveFlow() {
return IntegrationFlow
.from(() -> new GenericMessage<>("a,b,c,d,e,f"),
e -> e.poller(p -> p.trigger(ctx -> this.invoked.getAndSet(true) ? null : new Date()))
e -> e.poller(p -> p.trigger(ctx -> this.invoked.getAndSet(true) ? null : Instant.now()))
.id("reactiveStreamsMessageSource"))
.split(String.class, p -> p.split(","))
.log()

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2019 the original author or authors.
* Copyright 2002-2022 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.
@@ -19,13 +19,14 @@ package org.springframework.integration.endpoint;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.Mockito.mock;
import java.time.Duration;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.atomic.AtomicInteger;
import org.junit.Test;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.BeanFactory;
import org.springframework.expression.Expression;
@@ -40,39 +41,42 @@ import org.springframework.scheduling.support.PeriodicTrigger;
/**
* @author Mark Fisher
* @author Gary Russell
* @author Artem Bilan
*
* @since 2.0
*/
public class ExpressionEvaluatingMessageSourceIntegrationTests {
private static final AtomicInteger counter = new AtomicInteger();
@Test
public void test() throws Exception {
QueueChannel channel = new QueueChannel();
String payloadExpression = "'test-' + T(org.springframework.integration.endpoint.ExpressionEvaluatingMessageSourceIntegrationTests).next()";
String payloadExpression =
"'test-' + T(org.springframework.integration.endpoint.ExpressionEvaluatingMessageSourceIntegrationTests).next()";
ThreadPoolTaskScheduler scheduler = new ThreadPoolTaskScheduler();
scheduler.afterPropertiesSet();
Map<String, Expression> headerExpressions = new HashMap<String, Expression>();
Map<String, Expression> headerExpressions = new HashMap<>();
headerExpressions.put("foo", new LiteralExpression("x"));
headerExpressions.put("bar", new SpelExpressionParser().parseExpression("7 * 6"));
ExpressionFactoryBean factoryBean = new ExpressionFactoryBean(payloadExpression);
factoryBean.afterPropertiesSet();
Expression expression = factoryBean.getObject();
ExpressionEvaluatingMessageSource<Object> source = new ExpressionEvaluatingMessageSource<Object>(expression, Object.class);
ExpressionEvaluatingMessageSource<Object> source =
new ExpressionEvaluatingMessageSource<>(expression, Object.class);
source.setBeanFactory(mock(BeanFactory.class));
source.setHeaderExpressions(headerExpressions);
SourcePollingChannelAdapter adapter = new SourcePollingChannelAdapter();
adapter.setSource(source);
adapter.setTaskScheduler(scheduler);
adapter.setMaxMessagesPerPoll(3);
adapter.setTrigger(new PeriodicTrigger(60000));
adapter.setTrigger(new PeriodicTrigger(Duration.ofSeconds(60)));
adapter.setOutputChannel(channel);
adapter.setErrorHandler(t -> {
throw new IllegalStateException("unexpected exception in test", t);
});
adapter.start();
List<Message<?>> messages = new ArrayList<Message<?>>();
List<Message<?>> messages = new ArrayList<>();
for (int i = 0; i < 3; i++) {
messages.add(channel.receive(1000));
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2014-2020 the original author or authors.
* Copyright 2014-2022 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.
@@ -22,6 +22,7 @@ import static org.mockito.Mockito.atLeast;
import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.verify;
import java.time.Duration;
import java.util.ArrayList;
import java.util.Collections;
import java.util.LinkedList;
@@ -335,8 +336,8 @@ public class PollerAdviceTests {
SourcePollingChannelAdapter adapter = new SourcePollingChannelAdapter();
final CountDownLatch latch = new CountDownLatch(5);
final LinkedList<Object> overridePresent = new LinkedList<>();
final CompoundTrigger compoundTrigger = new CompoundTrigger(new PeriodicTrigger(10));
Trigger override = spy(new PeriodicTrigger(5));
final CompoundTrigger compoundTrigger = new CompoundTrigger(new PeriodicTrigger(Duration.ofMillis(10)));
Trigger override = spy(new PeriodicTrigger(Duration.ofMillis(5)));
final CompoundTriggerAdvice advice = new CompoundTriggerAdvice(compoundTrigger, override);
adapter.setSource(() -> {
synchronized (overridePresent) {
@@ -359,7 +360,7 @@ public class PollerAdviceTests {
synchronized (overridePresent) {
assertThat(overridePresent.subList(0, 5)).containsExactly(null, override, null, override, null);
}
verify(override, atLeast(2)).nextExecutionTime(any(TriggerContext.class));
verify(override, atLeast(2)).nextExecution(any(TriggerContext.class));
}
private void configure(SourcePollingChannelAdapter adapter) {

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2019 the original author or authors.
* Copyright 2002-2022 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.
@@ -16,6 +16,8 @@
package org.springframework.integration.endpoint;
import java.time.Duration;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.GenericMessage;
import org.springframework.scheduling.support.PeriodicTrigger;
@@ -24,11 +26,12 @@ import org.springframework.scheduling.support.PeriodicTrigger;
* @author Jonas Partner
* @author Gary Russell
* @author Andreas Baer
* @author Artem Bilan
*/
public class PollingEndpointStub extends AbstractPollingEndpoint {
public PollingEndpointStub() {
this.setTrigger(new PeriodicTrigger(500));
this.setTrigger(new PeriodicTrigger(Duration.ofMillis(500)));
}
@Override
@@ -38,7 +41,7 @@ public class PollingEndpointStub extends AbstractPollingEndpoint {
@Override
protected Message<?> receiveMessage() {
return new GenericMessage<String>("test message");
return new GenericMessage<>("test message");
}
@Override
@@ -50,4 +53,5 @@ public class PollingEndpointStub extends AbstractPollingEndpoint {
protected String getResourceKey() {
return "PollingEndpointStub";
}
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2021 the original author or authors.
* Copyright 2002-2022 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.
@@ -22,6 +22,7 @@ import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.times;
import java.time.Duration;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
@@ -85,7 +86,7 @@ public class PollingLifecycleTests {
});
PollingConsumer consumer = new PollingConsumer(channel, handler);
consumer.setTrigger(new PeriodicTrigger(0));
consumer.setTrigger(new PeriodicTrigger(Duration.ZERO));
consumer.setErrorHandler(errorHandler);
consumer.setTaskScheduler(taskScheduler);
consumer.setBeanFactory(mock(BeanFactory.class));
@@ -107,7 +108,7 @@ public class PollingLifecycleTests {
SourcePollingChannelAdapterFactoryBean adapterFactory = new SourcePollingChannelAdapterFactoryBean();
PollerMetadata pollerMetadata = new PollerMetadata();
pollerMetadata.setTrigger(new PeriodicTrigger(2000));
pollerMetadata.setTrigger(new PeriodicTrigger(Duration.ofSeconds(2)));
adapterFactory.setPollerMetadata(pollerMetadata);
//Has to be an explicit implementation - Mockito cannot mock/spy lambdas
@@ -140,7 +141,7 @@ public class PollingLifecycleTests {
SourcePollingChannelAdapterFactoryBean adapterFactory = new SourcePollingChannelAdapterFactoryBean();
PollerMetadata pollerMetadata = new PollerMetadata();
pollerMetadata.setMaxMessagesPerPoll(-1);
pollerMetadata.setTrigger(new PeriodicTrigger(2000));
pollerMetadata.setTrigger(new PeriodicTrigger(Duration.ofSeconds(2)));
adapterFactory.setPollerMetadata(pollerMetadata);
final Runnable caughtInterrupted = mock(Runnable.class);
final CountDownLatch interruptedLatch = new CountDownLatch(1);
@@ -180,7 +181,7 @@ public class PollingLifecycleTests {
adapterFactory.setOutputChannel(new NullChannel());
adapterFactory.setBeanFactory(mock(ConfigurableBeanFactory.class));
PollerMetadata pollerMetadata = new PollerMetadata();
pollerMetadata.setTrigger(new PeriodicTrigger(2000));
pollerMetadata.setTrigger(new PeriodicTrigger(Duration.ofSeconds(2)));
adapterFactory.setPollerMetadata(pollerMetadata);
final AtomicBoolean startInvoked = new AtomicBoolean();

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2016-2020 the original author or authors.
* Copyright 2016-2022 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.
@@ -20,6 +20,7 @@ import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatIllegalStateException;
import java.io.ByteArrayOutputStream;
import java.time.Duration;
import java.util.Arrays;
import java.util.HashMap;
import java.util.List;
@@ -334,7 +335,7 @@ public class IntegrationGraphServerTests {
@Bean(name = PollerMetadata.DEFAULT_POLLER)
public PollerMetadata defaultPoller() {
PollerMetadata poller = new PollerMetadata();
poller.setTrigger(new PeriodicTrigger(60000));
poller.setTrigger(new PeriodicTrigger(Duration.ofSeconds(60)));
MessagePublishingErrorHandler errorHandler = new MessagePublishingErrorHandler();
errorHandler.setDefaultErrorChannel(myErrors());
poller.setErrorHandler(errorHandler);

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2019 the original author or authors.
* Copyright 2002-2022 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.
@@ -17,13 +17,17 @@
package org.springframework.integration.message;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatExceptionOfType;
import static org.assertj.core.api.Assertions.assertThatIllegalStateException;
import static org.assertj.core.api.Assertions.assertThatNoException;
import static org.mockito.Mockito.mock;
import java.time.Duration;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.SynchronousQueue;
import java.util.concurrent.TimeUnit;
import org.junit.Test;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.BeanFactory;
import org.springframework.integration.channel.QueueChannel;
@@ -47,7 +51,8 @@ public class MethodInvokingMessageHandlerTests {
public void validMethod() {
MethodInvokingMessageHandler handler = new MethodInvokingMessageHandler(new TestSink(), "validMethod");
handler.setBeanFactory(mock(BeanFactory.class));
handler.handleMessage(new GenericMessage<>("test"));
assertThatNoException()
.isThrownBy(() -> handler.handleMessage(new GenericMessage<>("test")));
}
@Test
@@ -55,23 +60,21 @@ public class MethodInvokingMessageHandlerTests {
new MethodInvokingMessageHandler(new TestSink(), "validMethodWithNoArgs");
}
@Test(expected = MessagingException.class)
@Test
public void methodWithReturnValue() {
Message<?> message = new GenericMessage<>("test");
try {
MethodInvokingMessageHandler handler = new MethodInvokingMessageHandler(new TestSink(),
"methodWithReturnValue");
handler.handleMessage(message);
}
catch (MessagingException e) {
assertThat(message).isEqualTo(e.getFailedMessage());
throw e;
}
MethodInvokingMessageHandler handler = new MethodInvokingMessageHandler(new TestSink(),
"methodWithReturnValue");
assertThatExceptionOfType(MessagingException.class)
.isThrownBy(() -> handler.handleMessage(message))
.satisfies(messagingException ->
assertThat(messagingException.getFailedMessage()).isEqualTo(message));
}
@Test(expected = IllegalStateException.class)
@Test
public void noMatchingMethodName() {
new MethodInvokingMessageHandler(new TestSink(), "noSuchMethod");
assertThatIllegalStateException()
.isThrownBy(() -> new MethodInvokingMessageHandler(new TestSink(), "noSuchMethod"));
}
@Test
@@ -87,7 +90,7 @@ public class MethodInvokingMessageHandlerTests {
MethodInvokingMessageHandler handler = new MethodInvokingMessageHandler(testBean, "foo");
handler.setBeanFactory(context);
PollingConsumer endpoint = new PollingConsumer(channel, handler);
endpoint.setTrigger(new PeriodicTrigger(10));
endpoint.setTrigger(new PeriodicTrigger(Duration.ofMillis(10)));
context.registerEndpoint("testEndpoint", endpoint);
context.refresh();
String result = queue.poll(2000, TimeUnit.MILLISECONDS);
@@ -97,13 +100,7 @@ public class MethodInvokingMessageHandlerTests {
}
private static class TestBean {
private final BlockingQueue<String> queue;
TestBean(BlockingQueue<String> queue) {
this.queue = queue;
}
private record TestBean(BlockingQueue<String> queue) {
@SuppressWarnings("unused")
public void foo(String s) {

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2019 the original author or authors.
* Copyright 2002-2022 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.
@@ -18,14 +18,17 @@ package org.springframework.integration.transformer;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatExceptionOfType;
import static org.assertj.core.api.Assertions.assertThatIllegalArgumentException;
import static org.assertj.core.api.Assertions.assertThatIllegalStateException;
import static org.assertj.core.api.Assertions.fail;
import static org.mockito.Mockito.mock;
import java.time.Duration;
import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
import org.junit.Test;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.BeanFactory;
import org.springframework.beans.factory.BeanInitializationException;
@@ -74,7 +77,6 @@ public class ContentEnricherTests {
*/
@Test
public void replyChannelReplyTimingOut() throws Exception {
final long requestTimeout = 500L;
final long replyTimeout = 100L;
@@ -93,7 +95,7 @@ public class ContentEnricherTests {
expressionFactoryBean.setSingleton(false);
expressionFactoryBean.afterPropertiesSet();
final Map<String, Expression> expressions = new HashMap<String, Expression>();
final Map<String, Expression> expressions = new HashMap<>();
expressions.put("name", new LiteralExpression("cartman"));
expressions.put("child.name", expressionFactoryBean.getObject());
@@ -103,20 +105,21 @@ public class ContentEnricherTests {
enricher.setBeanFactory(mock(BeanFactory.class));
enricher.afterPropertiesSet();
final AbstractReplyProducingMessageHandler handler = new AbstractReplyProducingMessageHandler() {
final AbstractReplyProducingMessageHandler handler =
new AbstractReplyProducingMessageHandler() {
@Override
protected Object handleRequestMessage(Message<?> requestMessage) {
try {
Thread.sleep(5000);
}
catch (InterruptedException e) {
fail(e.getMessage());
}
return new Target("child");
}
@Override
protected Object handleRequestMessage(Message<?> requestMessage) {
try {
Thread.sleep(5000);
}
catch (InterruptedException e) {
fail(e.getMessage());
}
return new Target("child");
}
};
};
handler.setBeanFactory(mock(BeanFactory.class));
handler.afterPropertiesSet();
@@ -124,7 +127,7 @@ public class ContentEnricherTests {
final PollingConsumer consumer = new PollingConsumer(requestChannel, handler);
final TestErrorHandler errorHandler = new TestErrorHandler();
consumer.setTrigger(new PeriodicTrigger(0));
consumer.setTrigger(new PeriodicTrigger(Duration.ZERO));
consumer.setErrorHandler(errorHandler);
ThreadPoolTaskScheduler taskScheduler = new ThreadPoolTaskScheduler();
@@ -139,27 +142,16 @@ public class ContentEnricherTests {
final Target target = new Target("replace me");
Message<?> requestMessage = MessageBuilder.withPayload(target).setReplyChannel(replyChannel).build();
try {
enricher.handleMessage(requestMessage);
}
catch (ReplyRequiredException e) {
assertThat(e.getMessage())
.isEqualTo("No reply produced by handler 'Enricher', and its 'requiresReply' property is set to " +
"true.");
return;
}
finally {
consumer.stop();
taskScheduler.destroy();
}
fail("ReplyRequiredException expected.");
assertThatExceptionOfType(ReplyRequiredException.class)
.isThrownBy(() -> enricher.handleMessage(requestMessage))
.withMessage("No reply produced by handler 'Enricher', and its 'requiresReply' property is set to true.");
consumer.stop();
taskScheduler.destroy();
}
@Test
public void requestChannelSendTimingOut() {
final String requestChannelName = "Request_Channel";
final long requestTimeout = 200L;
@@ -214,57 +206,36 @@ public class ContentEnricherTests {
@Test
public void setReplyChannelWithoutRequestChannel() {
QueueChannel replyChannel = new QueueChannel();
ContentEnricher enricher = new ContentEnricher();
enricher.setReplyChannel(replyChannel);
enricher.setBeanFactory(mock(BeanFactory.class));
try {
enricher.afterPropertiesSet();
}
catch (IllegalStateException e) {
assertThat(e.getMessage())
.isEqualTo("If the replyChannel is set, then the requestChannel must not be null");
return;
}
fail("Expected an exception.");
assertThatIllegalStateException()
.isThrownBy(enricher::afterPropertiesSet)
.withMessage("If the replyChannel is set, then the requestChannel must not be null");
}
@Test
public void setNullReplyTimeout() {
ContentEnricher enricher = new ContentEnricher();
enricher.setBeanFactory(mock(BeanFactory.class));
try {
enricher.setReplyTimeout(null);
}
catch (IllegalArgumentException e) {
assertThat(e.getMessage()).isEqualTo("replyTimeout must not be null");
return;
}
fail("Expected an exception.");
assertThatIllegalArgumentException()
.isThrownBy(() -> enricher.setReplyTimeout(null))
.withMessage("replyTimeout must not be null");
}
@Test
public void setNullRequestTimeout() {
ContentEnricher enricher = new ContentEnricher();
enricher.setBeanFactory(mock(BeanFactory.class));
try {
enricher.setRequestTimeout(null);
}
catch (IllegalArgumentException e) {
assertThat(e.getMessage()).isEqualTo("requestTimeout must not be null");
return;
}
fail("Expected an exception.");
assertThatIllegalArgumentException()
.isThrownBy(() -> enricher.setRequestTimeout(null))
.withMessage("requestTimeout must not be null");
}
@Test
@@ -272,7 +243,7 @@ public class ContentEnricherTests {
QueueChannel replyChannel = new QueueChannel();
ContentEnricher enricher = new ContentEnricher();
SpelExpressionParser parser = new SpelExpressionParser();
Map<String, Expression> propertyExpressions = new HashMap<String, Expression>();
Map<String, Expression> propertyExpressions = new HashMap<>();
propertyExpressions.put("name", parser.parseExpression("'just a static string'"));
enricher.setPropertyExpressions(propertyExpressions);
enricher.setBeanFactory(mock(BeanFactory.class));
@@ -286,21 +257,13 @@ public class ContentEnricherTests {
@Test
public void testContentEnricherWithNullRequestChannel() {
ContentEnricher enricher = new ContentEnricher();
enricher.setReplyChannel(new QueueChannel());
enricher.setBeanFactory(mock(BeanFactory.class));
try {
enricher.afterPropertiesSet();
}
catch (IllegalStateException e) {
assertThat(e.getMessage())
.isEqualTo("If the replyChannel is set, then the requestChannel must not be null");
return;
}
fail("Expected an IllegalArgumentException to be thrown.");
assertThatIllegalStateException()
.isThrownBy(enricher::afterPropertiesSet)
.withMessage("If the replyChannel is set, then the requestChannel must not be null");
}
@Test
@@ -318,7 +281,7 @@ public class ContentEnricherTests {
enricher.setRequestChannel(requestChannel);
SpelExpressionParser parser = new SpelExpressionParser();
Map<String, Expression> propertyExpressions = new HashMap<String, Expression>();
Map<String, Expression> propertyExpressions = new HashMap<>();
propertyExpressions.put("child.name", parser.parseExpression("payload.lastName + ', ' + payload.firstName"));
enricher.setPropertyExpressions(propertyExpressions);
enricher.setBeanFactory(mock(BeanFactory.class));
@@ -349,7 +312,7 @@ public class ContentEnricherTests {
enricher.setShouldClonePayload(true);
SpelExpressionParser parser = new SpelExpressionParser();
Map<String, Expression> propertyExpressions = new HashMap<String, Expression>();
Map<String, Expression> propertyExpressions = new HashMap<>();
propertyExpressions.put("name", parser.parseExpression("payload.lastName + ', ' + payload.firstName"));
enricher.setPropertyExpressions(propertyExpressions);
enricher.setBeanFactory(mock(BeanFactory.class));
@@ -380,7 +343,7 @@ public class ContentEnricherTests {
enricher.setShouldClonePayload(true);
SpelExpressionParser parser = new SpelExpressionParser();
Map<String, Expression> propertyExpressions = new HashMap<String, Expression>();
Map<String, Expression> propertyExpressions = new HashMap<>();
propertyExpressions.put("name", parser.parseExpression("payload.lastName + ', ' + payload.firstName"));
enricher.setPropertyExpressions(propertyExpressions);
enricher.setBeanFactory(mock(BeanFactory.class));
@@ -414,7 +377,7 @@ public class ContentEnricherTests {
enricher.setShouldClonePayload(true);
SpelExpressionParser parser = new SpelExpressionParser();
Map<String, Expression> propertyExpressions = new HashMap<String, Expression>();
Map<String, Expression> propertyExpressions = new HashMap<>();
propertyExpressions.put("name", parser.parseExpression("payload.lastName + ', ' + payload.firstName"));
enricher.setPropertyExpressions(propertyExpressions);
enricher.setBeanFactory(mock(BeanFactory.class));
@@ -425,16 +388,10 @@ public class ContentEnricherTests {
Message<?> requestMessage = MessageBuilder.withPayload(target).setReplyChannel(replyChannel).build();
try {
enricher.handleMessage(requestMessage);
}
catch (MessageHandlingException e) {
assertThat(e.getMessage()).contains("Failed to clone payload object");
return;
}
fail("Expected a MessageHandlingException to be thrown.");
assertThatExceptionOfType(MessageHandlingException.class)
.isThrownBy(() -> enricher.handleMessage(requestMessage))
.withMessageContaining("Failed to clone payload object");
}
@Test
@@ -451,7 +408,6 @@ public class ContentEnricherTests {
@Test
public void testLifeCycleMethodsWithRequestChannel() {
DirectChannel requestChannel = new DirectChannel();
requestChannel.subscribe(new AbstractReplyProducingMessageHandler() {
@@ -484,9 +440,8 @@ public class ContentEnricherTests {
* the "error-channel" which returns a alternative {@link Target}.
*/
@Test
public void testErrorChannel() throws Exception {
final DirectChannel requestChannel = new DirectChannel();
public void testErrorChannel() {
DirectChannel requestChannel = new DirectChannel();
requestChannel.subscribe(new AbstractReplyProducingMessageHandler() {
@Override
@@ -496,7 +451,7 @@ public class ContentEnricherTests {
});
final DirectChannel errorChannel = new DirectChannel();
DirectChannel errorChannel = new DirectChannel();
errorChannel.subscribe(new AbstractReplyProducingMessageHandler() {
@Override
@@ -513,7 +468,7 @@ public class ContentEnricherTests {
enricher.setErrorChannel(errorChannel);
SpelExpressionParser parser = new SpelExpressionParser();
Map<String, Expression> propertyExpressions = new HashMap<String, Expression>();
Map<String, Expression> propertyExpressions = new HashMap<>();
propertyExpressions.put("name", parser.parseExpression("payload.name + ' target'"));
enricher.setPropertyExpressions(propertyExpressions);
@@ -536,15 +491,10 @@ public class ContentEnricherTests {
Collections.singletonMap(MessageHeaders.TIMESTAMP, new StaticHeaderValueMessageProcessor<>("foo")));
contentEnricher.setBeanFactory(mock(BeanFactory.class));
try {
contentEnricher.afterPropertiesSet();
fail("BeanInitializationException expected");
}
catch (Exception e) {
assertThat(e).isInstanceOf(BeanInitializationException.class);
assertThat(e.getMessage())
.contains("ContentEnricher cannot override 'id' and 'timestamp' read-only headers.");
}
assertThatExceptionOfType(BeanInitializationException.class)
.isThrownBy(contentEnricher::afterPropertiesSet)
.withMessageContaining("ContentEnricher cannot override 'id' and 'timestamp' read-only headers.");
}
@Test
@@ -554,28 +504,14 @@ public class ContentEnricherTests {
Collections.singletonMap(MessageHeaders.ID, new StaticHeaderValueMessageProcessor<>("foo")));
contentEnricher.setBeanFactory(mock(BeanFactory.class));
try {
contentEnricher.afterPropertiesSet();
fail("BeanInitializationException expected");
}
catch (Exception e) {
assertThat(e).isInstanceOf(BeanInitializationException.class);
assertThat(e.getMessage())
.contains("ContentEnricher cannot override 'id' and 'timestamp' read-only headers.");
}
assertThatExceptionOfType(BeanInitializationException.class)
.isThrownBy(contentEnricher::afterPropertiesSet)
.withMessageContaining("ContentEnricher cannot override 'id' and 'timestamp' read-only headers.");
}
@SuppressWarnings("unused")
private static final class Source {
private final String firstName;
private final String lastName;
Source(String firstName, String lastName) {
this.firstName = firstName;
this.lastName = lastName;
}
private record Source(String firstName, String lastName) {
public String getFirstName() {
return firstName;

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2016-2020 the original author or authors.
* Copyright 2016-2022 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.
@@ -20,6 +20,7 @@ import static org.assertj.core.api.Assertions.assertThat;
import java.io.Closeable;
import java.io.InputStream;
import java.time.Duration;
import java.util.Comparator;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.ConcurrentHashMap;
@@ -181,7 +182,7 @@ public class FtpStreamingMessageSourceTests extends FtpTestSupport {
@Bean(name = PollerMetadata.DEFAULT_POLLER)
public PollerMetadata defaultPoller() {
PollerMetadata pollerMetadata = new PollerMetadata();
pollerMetadata.setTrigger(new PeriodicTrigger(500));
pollerMetadata.setTrigger(new PeriodicTrigger(Duration.ofMillis(500)));
pollerMetadata.setMaxMessagesPerPoll(2);
return pollerMetadata;
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2020 the original author or authors.
* Copyright 2002-2022 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.
@@ -16,6 +16,7 @@
package org.springframework.integration.ip.tcp;
import java.time.Duration;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ScheduledFuture;
@@ -37,6 +38,7 @@ import org.springframework.integration.ip.tcp.connection.TcpSender;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessagingException;
import org.springframework.messaging.support.ErrorMessage;
import org.springframework.scheduling.TaskScheduler;
import org.springframework.util.Assert;
/**
@@ -91,7 +93,7 @@ public class TcpInboundGateway extends MessagingGatewaySupport implements
else {
if (isErrorMessage) {
/*
* Socket errors are sent here so they can be conveyed to any waiting thread.
* Socket errors are sent here, so they can be conveyed to any waiting thread.
* There's not one here; simply ignore.
*/
return false;
@@ -224,8 +226,9 @@ public class TcpInboundGateway extends MessagingGatewaySupport implements
ClientModeConnectionManager manager =
new ClientModeConnectionManager(this.clientConnectionFactory);
this.clientModeConnectionManager = manager;
Assert.state(getTaskScheduler() != null, "Client mode requires a task scheduler");
this.scheduledFuture = getTaskScheduler().scheduleAtFixedRate(manager, this.retryInterval);
TaskScheduler taskScheduler = getTaskScheduler();
Assert.state(taskScheduler != null, "Client mode requires a task scheduler");
this.scheduledFuture = taskScheduler.scheduleAtFixedRate(manager, Duration.ofMillis(this.retryInterval));
}
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2020 the original author or authors.
* Copyright 2002-2022 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.
@@ -19,6 +19,7 @@ package org.springframework.integration.ip.tcp.connection;
import static org.assertj.core.api.Assertions.assertThat;
import static org.awaitility.Awaitility.with;
import java.net.InetAddress;
import java.net.Socket;
import java.time.Duration;
import java.util.ArrayList;
@@ -48,18 +49,19 @@ import org.springframework.messaging.support.ErrorMessage;
*
* @since 2.0
*/
//@LongRunningTest
public class TcpNioConnectionReadTests {
private final CountDownLatch latch = new CountDownLatch(1);
private AbstractServerConnectionFactory getConnectionFactory(
AbstractByteArraySerializer serializer, TcpListener listener) throws Exception {
AbstractByteArraySerializer serializer, TcpListener listener) {
return getConnectionFactory(serializer, listener, null);
}
private AbstractServerConnectionFactory getConnectionFactory(
AbstractByteArraySerializer serializer, TcpListener listener, TcpSender sender) throws Exception {
AbstractByteArraySerializer serializer, TcpListener listener, TcpSender sender) {
TcpNioServerConnectionFactory scf = new TcpNioServerConnectionFactory(0);
scf.setUsingDirectBuffers(true);
scf.setApplicationEventPublisher(e -> {
@@ -78,7 +80,7 @@ public class TcpNioConnectionReadTests {
@Test
public void testReadLength() throws Exception {
ByteArrayLengthHeaderSerializer serializer = new ByteArrayLengthHeaderSerializer();
final List<Message<?>> responses = new ArrayList<Message<?>>();
final List<Message<?>> responses = new ArrayList<>();
final Semaphore semaphore = new Semaphore(0);
AbstractServerConnectionFactory scf = getConnectionFactory(serializer, message -> {
responses.add(message);
@@ -106,7 +108,7 @@ public class TcpNioConnectionReadTests {
@Test
public void testFragmented() throws Exception {
ByteArrayLengthHeaderSerializer serializer = new ByteArrayLengthHeaderSerializer();
final List<Message<?>> responses = new ArrayList<Message<?>>();
final List<Message<?>> responses = new ArrayList<>();
final Semaphore semaphore = new Semaphore(0);
AbstractServerConnectionFactory scf = getConnectionFactory(serializer, message -> {
responses.add(message);
@@ -137,7 +139,7 @@ public class TcpNioConnectionReadTests {
@Test
public void testReadStxEtx() throws Exception {
ByteArrayStxEtxSerializer serializer = new ByteArrayStxEtxSerializer();
final List<Message<?>> responses = new ArrayList<Message<?>>();
final List<Message<?>> responses = new ArrayList<>();
final Semaphore semaphore = new Semaphore(0);
AbstractServerConnectionFactory scf = getConnectionFactory(serializer, message -> {
responses.add(message);
@@ -164,7 +166,7 @@ public class TcpNioConnectionReadTests {
@Test
public void testReadCrLf() throws Exception {
ByteArrayCrLfSerializer serializer = new ByteArrayCrLfSerializer();
final List<Message<?>> responses = new ArrayList<Message<?>>();
final List<Message<?>> responses = new ArrayList<>();
final Semaphore semaphore = new Semaphore(0);
AbstractServerConnectionFactory scf = getConnectionFactory(serializer, message -> {
responses.add(message);
@@ -191,11 +193,11 @@ public class TcpNioConnectionReadTests {
public void testReadLengthOverflow() throws Exception {
ByteArrayLengthHeaderSerializer serializer = new ByteArrayLengthHeaderSerializer();
final Semaphore semaphore = new Semaphore(0);
final List<TcpConnection> added = new ArrayList<TcpConnection>();
final List<TcpConnection> removed = new ArrayList<TcpConnection>();
final List<TcpConnection> added = new ArrayList<>();
final List<TcpConnection> removed = new ArrayList<>();
final CountDownLatch errorMessageLetch = new CountDownLatch(1);
final AtomicReference<Throwable> errorMessageRef = new AtomicReference<Throwable>();
final AtomicReference<Throwable> errorMessageRef = new AtomicReference<>();
AbstractServerConnectionFactory scf = getConnectionFactory(serializer, message -> {
if (message instanceof ErrorMessage) {
@@ -299,11 +301,11 @@ public class TcpNioConnectionReadTests {
ByteArrayCrLfSerializer serializer = new ByteArrayCrLfSerializer();
serializer.setMaxMessageSize(1024);
final Semaphore semaphore = new Semaphore(0);
final List<TcpConnection> added = new ArrayList<TcpConnection>();
final List<TcpConnection> removed = new ArrayList<TcpConnection>();
final List<TcpConnection> added = new ArrayList<>();
final List<TcpConnection> removed = new ArrayList<>();
final CountDownLatch errorMessageLetch = new CountDownLatch(1);
final AtomicReference<Throwable> errorMessageRef = new AtomicReference<Throwable>();
final AtomicReference<Throwable> errorMessageRef = new AtomicReference<>();
AbstractServerConnectionFactory scf = getConnectionFactory(serializer, message -> {
if (message instanceof ErrorMessage) {
@@ -355,11 +357,11 @@ public class TcpNioConnectionReadTests {
ByteArrayCrLfSerializer serializer = new ByteArrayCrLfSerializer();
serializer.setMaxMessageSize(1024);
final Semaphore semaphore = new Semaphore(0);
final List<TcpConnection> added = new ArrayList<TcpConnection>();
final List<TcpConnection> removed = new ArrayList<TcpConnection>();
final List<TcpConnection> added = new ArrayList<>();
final List<TcpConnection> removed = new ArrayList<>();
final CountDownLatch errorMessageLetch = new CountDownLatch(1);
final AtomicReference<Throwable> errorMessageRef = new AtomicReference<Throwable>();
final AtomicReference<Throwable> errorMessageRef = new AtomicReference<>();
AbstractServerConnectionFactory scf = getConnectionFactory(serializer, message -> {
if (message instanceof ErrorMessage) {
@@ -408,11 +410,11 @@ public class TcpNioConnectionReadTests {
ByteArrayCrLfSerializer serializer = new ByteArrayCrLfSerializer();
serializer.setMaxMessageSize(1024);
final Semaphore semaphore = new Semaphore(0);
final List<TcpConnection> added = new ArrayList<TcpConnection>();
final List<TcpConnection> removed = new ArrayList<TcpConnection>();
final List<TcpConnection> added = new ArrayList<>();
final List<TcpConnection> removed = new ArrayList<>();
final CountDownLatch errorMessageLetch = new CountDownLatch(1);
final AtomicReference<Throwable> errorMessageRef = new AtomicReference<Throwable>();
final AtomicReference<Throwable> errorMessageRef = new AtomicReference<>();
AbstractServerConnectionFactory scf = getConnectionFactory(serializer, message -> {
if (message instanceof ErrorMessage) {
@@ -435,7 +437,7 @@ public class TcpNioConnectionReadTests {
}
});
Socket socket = SocketFactory.getDefault().createSocket("localhost", scf.getPort());
Socket socket = SocketFactory.getDefault().createSocket(InetAddress.getLocalHost(), scf.getPort());
socket.getOutputStream().write("partial".getBytes());
socket.close();
whileOpen(semaphore, added);
@@ -486,11 +488,11 @@ public class TcpNioConnectionReadTests {
private void testClosureMidMessageGuts(AbstractByteArraySerializer serializer, String shortMessage)
throws Exception {
final Semaphore semaphore = new Semaphore(0);
final List<TcpConnection> added = new ArrayList<TcpConnection>();
final List<TcpConnection> removed = new ArrayList<TcpConnection>();
final List<TcpConnection> added = new ArrayList<>();
final List<TcpConnection> removed = new ArrayList<>();
final CountDownLatch errorMessageLetch = new CountDownLatch(1);
final AtomicReference<Throwable> errorMessageRef = new AtomicReference<Throwable>();
final AtomicReference<Throwable> errorMessageRef = new AtomicReference<>();
AbstractServerConnectionFactory scf = getConnectionFactory(serializer, message -> {
if (message instanceof ErrorMessage) {
@@ -513,7 +515,7 @@ public class TcpNioConnectionReadTests {
}
});
Socket socket = SocketFactory.getDefault().createSocket("localhost", scf.getPort());
Socket socket = SocketFactory.getDefault().createSocket(InetAddress.getLocalHost(), scf.getPort());
socket.getOutputStream().write(shortMessage.getBytes());
socket.close();
whileOpen(semaphore, added);

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2020 the original author or authors.
* Copyright 2002-2022 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.
@@ -62,7 +62,7 @@ public class SocketTestUtils {
public static CountDownLatch testSendLength(final int port, final CountDownLatch latch) {
final CountDownLatch testCompleteLatch = new CountDownLatch(1);
Thread thread = new Thread(() -> {
try (Socket socket = new Socket(InetAddress.getByName("localhost"), port)) {
try (Socket socket = new Socket(InetAddress.getLocalHost(), port)) {
for (int i = 0; i < 2; i++) {
byte[] len = new byte[4];
ByteBuffer.wrap(len).putInt(TEST_STRING.length() * 2);
@@ -94,7 +94,7 @@ public class SocketTestUtils {
public static CountDownLatch testSendLengthOverflow(final int port) {
final CountDownLatch testCompleteLatch = new CountDownLatch(1);
Thread thread = new Thread(() -> {
try (Socket socket = new Socket(InetAddress.getByName("localhost"), port)) {
try (Socket socket = new Socket(InetAddress.getLocalHost(), port)) {
byte[] len = new byte[4];
ByteBuffer.wrap(len).putInt(Integer.MAX_VALUE);
socket.getOutputStream().write(len);
@@ -120,7 +120,7 @@ public class SocketTestUtils {
Socket socket = null;
try {
logger.debug("Connecting to " + port);
socket = new Socket(InetAddress.getByName("localhost"), port);
socket = new Socket(InetAddress.getLocalHost(), port);
OutputStream os = socket.getOutputStream();
for (int i = 0; i < howMany; i++) {
writeByte(os, 0, noDelay);
@@ -162,12 +162,12 @@ public class SocketTestUtils {
/**
* Sends a STX/ETX message in two chunks. Two such messages are sent.
* @param latch If not null, await until counted down before sending second chunk.
* @param latch If not null, waits until counted down before sending second chunk.
*/
public static CountDownLatch testSendStxEtx(final int port, final CountDownLatch latch) {
final CountDownLatch testCompleteLatch = new CountDownLatch(1);
Thread thread = new Thread(() -> {
try (Socket socket = new Socket(InetAddress.getByName("localhost"), port)) {
try (Socket socket = new Socket(InetAddress.getLocalHost(), port)) {
OutputStream outputStream = socket.getOutputStream();
for (int i = 0; i < 2; i++) {
writeByte(outputStream, 0x02, true);
@@ -199,7 +199,7 @@ public class SocketTestUtils {
public static CountDownLatch testSendStxEtxOverflow(final int port) {
final CountDownLatch testCompleteLatch = new CountDownLatch(1);
Thread thread = new Thread(() -> {
try (Socket socket = new Socket(InetAddress.getByName("localhost"), port)) {
try (Socket socket = new Socket(InetAddress.getLocalHost(), port)) {
OutputStream outputStream = socket.getOutputStream();
writeByte(outputStream, 0x02, true);
for (int i = 0; i < 1500; i++) {
@@ -223,7 +223,7 @@ public class SocketTestUtils {
public static CountDownLatch testSendCrLf(final int port, final CountDownLatch latch) {
final CountDownLatch testCompleteLatch = new CountDownLatch(1);
Thread thread = new Thread(() -> {
try (Socket socket = new Socket(InetAddress.getByName("localhost"), port)) {
try (Socket socket = new Socket(InetAddress.getLocalHost(), port)) {
OutputStream outputStream = socket.getOutputStream();
for (int i = 0; i < 2; i++) {
outputStream.write(TEST_STRING.getBytes());
@@ -255,7 +255,7 @@ public class SocketTestUtils {
*/
public static void testSendCrLfSingle(final int port, final CountDownLatch latch) {
Thread thread = new Thread(() -> {
try (Socket socket = new Socket(InetAddress.getByName("localhost"), port)) {
try (Socket socket = new Socket(InetAddress.getLocalHost(), port)) {
OutputStream outputStream = socket.getOutputStream();
outputStream.write(TEST_STRING.getBytes());
outputStream.write(TEST_STRING.getBytes());
@@ -278,7 +278,7 @@ public class SocketTestUtils {
*/
public static void testSendRaw(final int port) {
Thread thread = new Thread(() -> {
try (Socket socket = new Socket(InetAddress.getByName("localhost"), port)) {
try (Socket socket = new Socket(InetAddress.getLocalHost(), port)) {
OutputStream outputStream = socket.getOutputStream();
outputStream.write(TEST_STRING.getBytes());
outputStream.write(TEST_STRING.getBytes());
@@ -298,7 +298,7 @@ public class SocketTestUtils {
public static CountDownLatch testSendSerialized(final int port) {
final CountDownLatch testCompleteLatch = new CountDownLatch(1);
Thread thread = new Thread(() -> {
try (Socket socket = new Socket(InetAddress.getByName("localhost"), port)) {
try (Socket socket = new Socket(InetAddress.getLocalHost(), port)) {
OutputStream outputStream = socket.getOutputStream();
ObjectOutputStream oos = new ObjectOutputStream(outputStream);
oos.writeObject(TEST_STRING);
@@ -323,7 +323,7 @@ public class SocketTestUtils {
public static CountDownLatch testSendCrLfOverflow(final int port) {
final CountDownLatch testCompleteLatch = new CountDownLatch(1);
Thread thread = new Thread(() -> {
try (Socket socket = new Socket(InetAddress.getByName("localhost"), port)) {
try (Socket socket = new Socket(InetAddress.getLocalHost(), port)) {
OutputStream outputStream = socket.getOutputStream();
for (int i = 0; i < 1500; i++) {
writeByte(outputStream, 'x', true);

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2021 the original author or authors.
* Copyright 2002-2022 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.
@@ -16,7 +16,8 @@
package org.springframework.integration.mail;
import java.util.Date;
import java.io.Serial;
import java.time.Instant;
import java.util.List;
import java.util.concurrent.Executor;
import java.util.concurrent.ExecutorService;
@@ -224,12 +225,12 @@ public class ImapIdleChannelAdapter extends MessageProducerSupport implements Be
};
}
private void publishException(Exception e) {
private void publishException(Exception ex) {
if (this.applicationEventPublisher != null) {
this.applicationEventPublisher.publishEvent(new ImapIdleExceptionEvent(e));
this.applicationEventPublisher.publishEvent(new ImapIdleExceptionEvent(ex));
}
else {
logger.debug(() -> "No application event publisher for exception: " + e.getMessage());
logger.debug(() -> "No application event publisher for exception: " + ex.getMessage());
}
}
@@ -305,13 +306,13 @@ public class ImapIdleChannelAdapter extends MessageProducerSupport implements Be
}
@Override
public Date nextExecutionTime(TriggerContext triggerContext) {
public Instant nextExecution(TriggerContext triggerContext) {
if (this.delayNextExecution) {
this.delayNextExecution = false;
return new Date(System.currentTimeMillis() + ImapIdleChannelAdapter.this.reconnectDelay);
return Instant.now().plusMillis(ImapIdleChannelAdapter.this.reconnectDelay);
}
else {
return new Date(System.currentTimeMillis());
return Instant.now();
}
}
@@ -323,10 +324,11 @@ public class ImapIdleChannelAdapter extends MessageProducerSupport implements Be
public class ImapIdleExceptionEvent extends MailIntegrationEvent {
@Serial
private static final long serialVersionUID = -5875388810251967741L;
ImapIdleExceptionEvent(Exception e) {
super(ImapIdleChannelAdapter.this, e);
ImapIdleExceptionEvent(Exception ex) {
super(ImapIdleChannelAdapter.this, ex);
}
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2021 the original author or authors.
* Copyright 2002-2022 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.
@@ -16,8 +16,8 @@
package org.springframework.integration.mail;
import java.time.Instant;
import java.util.Arrays;
import java.util.Date;
import java.util.Objects;
import java.util.Properties;
import java.util.concurrent.ScheduledFuture;
@@ -185,8 +185,8 @@ public class ImapMailReceiver extends AbstractMailReceiver {
return;
}
try {
this.pingTask = this.scheduler.schedule(this.idleCanceler,
new Date(System.currentTimeMillis() + this.cancelIdleInterval));
this.pingTask =
this.scheduler.schedule(this.idleCanceler, Instant.now().plusMillis(this.cancelIdleInterval));
if (imapFolder.isOpen()) {
imapFolder.idle(true);
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2021 the original author or authors.
* Copyright 2002-2022 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.
@@ -16,8 +16,8 @@
package org.springframework.integration.mqtt.inbound;
import java.time.Instant;
import java.util.Arrays;
import java.util.Date;
import java.util.concurrent.ScheduledFuture;
import org.eclipse.paho.client.mqttv3.IMqttClient;
@@ -336,21 +336,22 @@ public class MqttPahoMessageDrivenChannelAdapter extends AbstractMqttMessageDriv
cancelReconnect();
if (isActive()) {
try {
this.reconnectFuture = getTaskScheduler().schedule(() -> {
try {
logger.debug("Attempting reconnect");
synchronized (MqttPahoMessageDrivenChannelAdapter.this) {
if (!MqttPahoMessageDrivenChannelAdapter.this.connected) {
connectAndSubscribe();
MqttPahoMessageDrivenChannelAdapter.this.reconnectFuture = null;
this.reconnectFuture = getTaskScheduler()
.schedule(() -> {
try {
logger.debug("Attempting reconnect");
synchronized (MqttPahoMessageDrivenChannelAdapter.this) {
if (!MqttPahoMessageDrivenChannelAdapter.this.connected) {
connectAndSubscribe();
MqttPahoMessageDrivenChannelAdapter.this.reconnectFuture = null;
}
}
}
}
}
catch (MqttException ex) {
logger.error(ex, "Exception while connecting and subscribing");
scheduleReconnect();
}
}, new Date(System.currentTimeMillis() + this.recoveryInterval));
catch (MqttException ex) {
logger.error(ex, "Exception while connecting and subscribing");
scheduleReconnect();
}
}, Instant.now().plusMillis(this.recoveryInterval));
}
catch (Exception ex) {
logger.error(ex, "Failed to schedule reconnect");

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2021 the original author or authors.
* Copyright 2002-2022 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.
@@ -37,7 +37,7 @@ import static org.mockito.Mockito.verify;
import java.lang.reflect.InvocationTargetException;
import java.lang.reflect.Method;
import java.util.Date;
import java.time.Instant;
import java.util.Properties;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.CountDownLatch;
@@ -402,7 +402,7 @@ public class MqttAdapterTests {
TaskScheduler taskScheduler = TestUtils.getPropertyValue(adapter, "taskScheduler", TaskScheduler.class);
verify(taskScheduler, never())
.schedule(any(Runnable.class), any(Date.class));
.schedule(any(Runnable.class), any(Instant.class));
}
@Test

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2016-2020 the original author or authors.
* Copyright 2016-2022 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.
@@ -19,6 +19,7 @@ package org.springframework.integration.sftp.inbound;
import static org.assertj.core.api.Assertions.assertThat;
import java.io.InputStream;
import java.time.Duration;
import java.util.Arrays;
import java.util.Comparator;
import java.util.concurrent.ConcurrentHashMap;
@@ -189,7 +190,7 @@ public class SftpStreamingMessageSourceTests extends SftpTestSupport {
@Bean(name = PollerMetadata.DEFAULT_POLLER)
public PollerMetadata defaultPoller() {
PollerMetadata pollerMetadata = new PollerMetadata();
pollerMetadata.setTrigger(new PeriodicTrigger(500));
pollerMetadata.setTrigger(new PeriodicTrigger(Duration.ofMillis(500)));
pollerMetadata.setMaxMessagesPerPoll(2000);
return pollerMetadata;
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2015-2020 the original author or authors.
* Copyright 2015-2022 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.
@@ -16,9 +16,9 @@
package org.springframework.integration.stomp;
import java.time.Instant;
import java.util.ArrayList;
import java.util.Collections;
import java.util.Date;
import java.util.List;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ScheduledFuture;
@@ -262,8 +262,7 @@ public abstract class AbstractStompSessionManager implements StompSessionManager
TaskScheduler taskScheduler = this.stompClient.getTaskScheduler();
if (taskScheduler != null) {
this.reconnectFuture =
taskScheduler.schedule(this::connect,
new Date(System.currentTimeMillis() + this.recoveryInterval));
taskScheduler.schedule(this::connect, Instant.now().plusMillis(this.recoveryInterval));
}
else {
this.logger.info("For automatic reconnection the stompClient should be configured with a TaskScheduler.");
@@ -278,7 +277,7 @@ public abstract class AbstractStompSessionManager implements StompSessionManager
this.reconnectFuture = null;
}
this.stompSessionListenableFuture.addCallback(
new ListenableFutureCallback<StompSession>() {
new ListenableFutureCallback<>() {
@Override
public void onFailure(Throwable ex) {

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2019 the original author or authors.
* Copyright 2002-2022 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.
@@ -22,10 +22,11 @@ import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
import java.io.StringReader;
import java.time.Duration;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import org.junit.Test;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.BeanFactory;
import org.springframework.context.ApplicationEventPublisher;
@@ -38,6 +39,7 @@ import org.springframework.scheduling.support.PeriodicTrigger;
/**
* @author Mark Fisher
* @author Gary Russell
* @author Artem Bilan
*/
public class CharacterStreamSourceTests {
@@ -85,7 +87,7 @@ public class CharacterStreamSourceTests {
ThreadPoolTaskScheduler scheduler = new ThreadPoolTaskScheduler();
scheduler.afterPropertiesSet();
adapter.setTaskScheduler(scheduler);
adapter.setTrigger(new PeriodicTrigger(100));
adapter.setTrigger(new PeriodicTrigger(Duration.ofMillis(100)));
QueueChannel out = new QueueChannel();
adapter.setOutputChannel(out);
adapter.setBeanFactory(mock(BeanFactory.class));

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2019 the original author or authors.
* Copyright 2002-2022 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.
@@ -20,14 +20,14 @@ import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.Mockito.mock;
import java.io.StringWriter;
import java.util.Date;
import java.time.Instant;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import org.junit.After;
import org.junit.Before;
import org.junit.Test;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.BeanFactory;
import org.springframework.integration.channel.QueueChannel;
@@ -39,6 +39,7 @@ import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
/**
* @author Mark Fisher
* @author Artem Bilan
*/
public class CharacterStreamWritingMessageHandlerTests {
@@ -55,7 +56,7 @@ public class CharacterStreamWritingMessageHandlerTests {
private ThreadPoolTaskScheduler scheduler;
@Before
@BeforeEach
public void initialize() {
writer = new StringWriter();
handler = new CharacterStreamWritingMessageHandler(writer);
@@ -70,15 +71,15 @@ public class CharacterStreamWritingMessageHandlerTests {
endpoint.setBeanFactory(mock(BeanFactory.class));
}
@After
public void stop() throws Exception {
@AfterEach
public void stop() {
scheduler.destroy();
}
@Test
public void singleString() {
handler.handleMessage(new GenericMessage<String>("foo"));
handler.handleMessage(new GenericMessage<>("foo"));
assertThat(writer.toString()).isEqualTo("foo");
}
@@ -86,8 +87,8 @@ public class CharacterStreamWritingMessageHandlerTests {
public void twoStringsAndNoNewLinesByDefault() {
endpoint.setMaxMessagesPerPoll(1);
endpoint.setTrigger(trigger);
channel.send(new GenericMessage<String>("foo"), 0);
channel.send(new GenericMessage<String>("bar"), 0);
channel.send(new GenericMessage<>("foo"), 0);
channel.send(new GenericMessage<>("bar"), 0);
endpoint.start();
trigger.await();
endpoint.stop();
@@ -104,8 +105,8 @@ public class CharacterStreamWritingMessageHandlerTests {
handler.setShouldAppendNewLine(true);
endpoint.setTrigger(trigger);
endpoint.setMaxMessagesPerPoll(1);
channel.send(new GenericMessage<String>("foo"), 0);
channel.send(new GenericMessage<String>("bar"), 0);
channel.send(new GenericMessage<>("foo"), 0);
channel.send(new GenericMessage<>("bar"), 0);
endpoint.start();
trigger.await();
endpoint.stop();
@@ -122,8 +123,8 @@ public class CharacterStreamWritingMessageHandlerTests {
public void maxMessagesPerTaskSameAsMessageCount() {
endpoint.setTrigger(trigger);
endpoint.setMaxMessagesPerPoll(2);
channel.send(new GenericMessage<String>("foo"), 0);
channel.send(new GenericMessage<String>("bar"), 0);
channel.send(new GenericMessage<>("foo"), 0);
channel.send(new GenericMessage<>("bar"), 0);
endpoint.start();
trigger.await();
endpoint.stop();
@@ -136,8 +137,8 @@ public class CharacterStreamWritingMessageHandlerTests {
endpoint.setMaxMessagesPerPoll(10);
endpoint.setReceiveTimeout(0);
handler.setShouldAppendNewLine(true);
channel.send(new GenericMessage<String>("foo"), 0);
channel.send(new GenericMessage<String>("bar"), 0);
channel.send(new GenericMessage<>("foo"), 0);
channel.send(new GenericMessage<>("bar"), 0);
endpoint.start();
trigger.await();
endpoint.stop();
@@ -150,7 +151,7 @@ public class CharacterStreamWritingMessageHandlerTests {
endpoint.setTrigger(trigger);
endpoint.setMaxMessagesPerPoll(1);
TestObject testObject = new TestObject("foo");
channel.send(new GenericMessage<TestObject>(testObject));
channel.send(new GenericMessage<>(testObject));
endpoint.start();
trigger.await();
endpoint.stop();
@@ -164,8 +165,8 @@ public class CharacterStreamWritingMessageHandlerTests {
endpoint.setMaxMessagesPerPoll(2);
TestObject testObject1 = new TestObject("foo");
TestObject testObject2 = new TestObject("bar");
channel.send(new GenericMessage<TestObject>(testObject1), 0);
channel.send(new GenericMessage<TestObject>(testObject2), 0);
channel.send(new GenericMessage<>(testObject1), 0);
channel.send(new GenericMessage<>(testObject2), 0);
endpoint.start();
trigger.await();
endpoint.stop();
@@ -180,8 +181,8 @@ public class CharacterStreamWritingMessageHandlerTests {
endpoint.setTrigger(trigger);
TestObject testObject1 = new TestObject("foo");
TestObject testObject2 = new TestObject("bar");
channel.send(new GenericMessage<TestObject>(testObject1), 0);
channel.send(new GenericMessage<TestObject>(testObject2), 0);
channel.send(new GenericMessage<>(testObject1), 0);
channel.send(new GenericMessage<>(testObject2), 0);
endpoint.start();
trigger.await();
endpoint.stop();
@@ -190,13 +191,7 @@ public class CharacterStreamWritingMessageHandlerTests {
}
private static class TestObject {
private final String text;
TestObject(String text) {
this.text = text;
}
private record TestObject(String text) {
@Override
public String toString() {
@@ -217,9 +212,9 @@ public class CharacterStreamWritingMessageHandlerTests {
}
@Override
public Date nextExecutionTime(TriggerContext triggerContext) {
public Instant nextExecution(TriggerContext triggerContext) {
if (!hasRun.getAndSet(true)) {
return new Date();
return Instant.now();
}
this.latch.countDown();
return null;

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2019 the original author or authors.
* Copyright 2002-2022 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.
@@ -16,7 +16,7 @@
package org.springframework.integration.test.util;
import java.util.Date;
import java.time.Instant;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
@@ -37,23 +37,23 @@ public class OnlyOnceTrigger implements Trigger {
private final AtomicBoolean hasRun = new AtomicBoolean();
private final Date executionTime;
private final Instant executionTime;
private volatile CountDownLatch latch = new CountDownLatch(1);
public OnlyOnceTrigger() {
this.executionTime = new Date();
this.executionTime = Instant.now();
}
@Override
public Date nextExecutionTime(TriggerContext triggerContext) {
public Instant nextExecution(TriggerContext triggerContext) {
if (this.hasRun.getAndSet(true)) {
this.latch.countDown();
return null;
}
return this.executionTime; // NOSONAR - expose internals
return this.executionTime;
}
@Override
@@ -90,9 +90,9 @@ public class OnlyOnceTrigger implements Trigger {
throw new IllegalStateException("test latch.await() did not count down");
}
}
catch (InterruptedException e) {
catch (InterruptedException ex) {
Thread.currentThread().interrupt();
throw new IllegalStateException("test latch.await() interrupted", e);
throw new IllegalStateException("test latch.await() interrupted", ex);
}
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2020 the original author or authors.
* Copyright 2002-2022 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 static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.verify;
import java.time.Duration;
import org.junit.jupiter.api.Test;
import org.mockito.ArgumentMatchers;
@@ -263,8 +265,9 @@ public class WebServiceOutboundGatewayParserTests {
assertThat(triggerObject.getClass()).isEqualTo(PeriodicTrigger.class);
PeriodicTrigger trigger = (PeriodicTrigger) triggerObject;
DirectFieldAccessor accessor = new DirectFieldAccessor(trigger);
assertThat(((Long) accessor.getPropertyValue("period")).longValue()).as("PeriodicTrigger had wrong period")
.isEqualTo(5000);
assertThat(accessor.getPropertyValue("period"))
.as("PeriodicTrigger had wrong period")
.isEqualTo(Duration.ofSeconds(5));
}
@Test

View File

@@ -1003,7 +1003,7 @@ catch { ... }
====
To improve performance, you may wish to send multiple messages and wait for the confirmations later, rather than one-at-a-time.
The returned message is the raw message after conversion; you can subclass `CorrelationData` with whatever additional data you need.
The returned message is the raw message after conversion; you can sub-class a `CorrelationData` with whatever additional data you need.
[[amqp-conversion-inbound]]
=== Inbound Message Conversion
@@ -1015,7 +1015,7 @@ If a conversion error occurs, and there is no error channel defined, the excepti
The default error handler treats conversion errors as fatal and the message will be rejected (and routed to a dead-letter exchange, if the queue is so configured).
If an error channel is defined, the `ErrorMessage` payload is a `ListenerExecutionFailedException` with properties `failedMessage` (the Spring AMQP message that could not be converted) and the `cause`.
If the container `AcknowledgeMode` is `AUTO` (the default) and the error flow consumes the error without throwing an exception, the original message will be acknowledged.
If the error flow throws an exception, the exception type, in conjunction with the container's error handler, will determine whether or not the message is requeued.
If the error flow throws an exception, the exception type, in conjunction with the container's error handler, will determine whether the message is requeued.
If the container is configured with `AcknowledgeMode.MANUAL`, the payload is a `ManualAckListenerExecutionFailedException` with additional properties `channel` and `deliveryTag`.
This enables the error flow to call `basicAck` or `basicNack` (or `basicReject`) for the message, to control its disposition.

View File

@@ -95,18 +95,18 @@ The following example shows how to set the trigger:
----
PollingConsumer consumer = new PollingConsumer(channel, handler);
consumer.setTrigger(new PeriodicTrigger(30, TimeUnit.SECONDS));
consumer.setTrigger(new PeriodicTrigger(Duration.ofSeconds(30)));
----
====
The `PeriodicTrigger` is typically defined with a simple interval (in milliseconds) but also supports an `initialDelay` property and a boolean `fixedRate` property (the default is `false` -- that is, no fixed delay).
The `PeriodicTrigger` is typically defined with a simple interval (`Duration`) but also supports an `initialDelay` property and a boolean `fixedRate` property (the default is `false` -- that is, no fixed delay).
The following example sets both properties:
====
[source,java]
----
PeriodicTrigger trigger = new PeriodicTrigger(1000);
trigger.setInitialDelay(5000);
PeriodicTrigger trigger = new PeriodicTrigger(Duration.ofSeconds(1));
trigger.setInitialDelay(Duration.ofSeconds(5));
trigger.setFixedRate(true);
----
====