diff --git a/build.gradle b/build.gradle index f0e33254d5..6c1a89ff55 100644 --- a/build.gradle +++ b/build.gradle @@ -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" diff --git a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java index fa6184a7f3..6249193ec3 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java @@ -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); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/DefaultHeaderChannelRegistry.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/DefaultHeaderChannelRegistry.java index eece643b8e..687c48f9d7 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/channel/DefaultHeaderChannelRegistry.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/DefaultHeaderChannelRegistry.java @@ -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; } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java index a1bdb36693..907e538d67 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java @@ -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 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 .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); } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java index 1b72b2d6d8..071c5a56cd 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java @@ -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 diff --git a/spring-integration-core/src/main/java/org/springframework/integration/util/CompoundTrigger.java b/spring-integration-core/src/main/java/org/springframework/integration/util/CompoundTrigger.java index 49a6e5f238..1956c15b1e 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/util/CompoundTrigger.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/util/CompoundTrigger.java @@ -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); } } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/util/DynamicPeriodicTrigger.java b/spring-integration-core/src/main/java/org/springframework/integration/util/DynamicPeriodicTrigger.java index 95ce697104..b152c10440 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/util/DynamicPeriodicTrigger.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/util/DynamicPeriodicTrigger.java @@ -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); } } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/bus/ApplicationContextMessageBusTests.java b/spring-integration-core/src/test/java/org/springframework/integration/bus/ApplicationContextMessageBusTests.java index efc85f150e..d4029b152e 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/bus/ApplicationContextMessageBusTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/bus/ApplicationContextMessageBusTests.java @@ -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 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("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 { - - private final CountDownLatch latch; - - FailingSource(CountDownLatch latch) { - this.latch = latch; - } + private record FailingSource(CountDownLatch latch) implements MessageSource { @Override public Message receive() { diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/SourcePollingChannelAdapterFactoryBeanTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/SourcePollingChannelAdapterFactoryBeanTests.java index 97a81010f1..de626618be 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/SourcePollingChannelAdapterFactoryBeanTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/SourcePollingChannelAdapterFactoryBeanTests.java @@ -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); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/StubTaskScheduler.java b/spring-integration-core/src/test/java/org/springframework/integration/config/StubTaskScheduler.java index 5dfa95d19f..8caed29947 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/StubTaskScheduler.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/StubTaskScheduler.java @@ -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; } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/TestTrigger.java b/spring-integration-core/src/test/java/org/springframework/integration/config/TestTrigger.java index 516368af61..122f68a62b 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/TestTrigger.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/TestTrigger.java @@ -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(); } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/InboundChannelAdapterExpressionTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/InboundChannelAdapterExpressionTests.java index ca7c472f34..8c19442842 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/InboundChannelAdapterExpressionTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/InboundChannelAdapterExpressionTests.java @@ -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")); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/InboundChannelAdapterWithDefaultPollerTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/InboundChannelAdapterWithDefaultPollerTests.java index eaa7c2a6f3..ce4da69679 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/InboundChannelAdapterWithDefaultPollerTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/InboundChannelAdapterWithDefaultPollerTests.java @@ -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(); } } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/IntervalTriggerParserTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/IntervalTriggerParserTests.java index 604eaba171..5ac3da0907 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/IntervalTriggerParserTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/IntervalTriggerParserTests.java @@ -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(); } + } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/PollerParserTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/PollerParserTests.java index 6a9edd2529..66d423bbef 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/PollerParserTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/PollerParserTests.java @@ -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)); } } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java b/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java index 072c3c061f..e94b01da99 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java @@ -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); } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/dsl/PollersTests.java b/spring-integration-core/src/test/java/org/springframework/integration/dsl/PollersTests.java index 45278b91d8..997c6a6b27 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/dsl/PollersTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/dsl/PollersTests.java @@ -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(); } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/dsl/flowservices/FlowServiceTests.java b/spring-integration-core/src/test/java/org/springframework/integration/dsl/flowservices/FlowServiceTests.java index c3ee4a3d75..cd0b5dcfe8 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/dsl/flowservices/FlowServiceTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/dsl/flowservices/FlowServiceTests.java @@ -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 executionDate = new AtomicReference<>(new Date()); + private final AtomicReference 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)) diff --git a/spring-integration-core/src/test/java/org/springframework/integration/dsl/manualflow/ManualFlowTests.java b/spring-integration-core/src/test/java/org/springframework/integration/dsl/manualflow/ManualFlowTests.java index 3ae58e2448..84451d9c4f 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/dsl/manualflow/ManualFlowTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/dsl/manualflow/ManualFlowTests.java @@ -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 nextExecutionTime = new AtomicReference<>(new Date()); + private final AtomicReference nextExecutionTime = new AtomicReference<>(Instant.now()); @Override protected IntegrationFlowDefinition buildFlow() { diff --git a/spring-integration-core/src/test/java/org/springframework/integration/dsl/reactivestreams/ReactiveStreamsTests.java b/spring-integration-core/src/test/java/org/springframework/integration/dsl/reactivestreams/ReactiveStreamsTests.java index a6b59a24e8..7fdcf60eb1 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/dsl/reactivestreams/ReactiveStreamsTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/dsl/reactivestreams/ReactiveStreamsTests.java @@ -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> 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() diff --git a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/ExpressionEvaluatingMessageSourceIntegrationTests.java b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/ExpressionEvaluatingMessageSourceIntegrationTests.java index 439cbe317d..5025f704d7 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/ExpressionEvaluatingMessageSourceIntegrationTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/ExpressionEvaluatingMessageSourceIntegrationTests.java @@ -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 headerExpressions = new HashMap(); + Map 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 source = new ExpressionEvaluatingMessageSource(expression, Object.class); + ExpressionEvaluatingMessageSource 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> messages = new ArrayList>(); + List> messages = new ArrayList<>(); for (int i = 0; i < 3; i++) { messages.add(channel.receive(1000)); } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollerAdviceTests.java b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollerAdviceTests.java index 5e08c409bf..384070d295 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollerAdviceTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollerAdviceTests.java @@ -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 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) { diff --git a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingEndpointStub.java b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingEndpointStub.java index e660e35cf5..3f91e4e07c 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingEndpointStub.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingEndpointStub.java @@ -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("test message"); + return new GenericMessage<>("test message"); } @Override @@ -50,4 +53,5 @@ public class PollingEndpointStub extends AbstractPollingEndpoint { protected String getResourceKey() { return "PollingEndpointStub"; } + } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingLifecycleTests.java b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingLifecycleTests.java index 3951571546..a27488d07d 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingLifecycleTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingLifecycleTests.java @@ -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(); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/graph/IntegrationGraphServerTests.java b/spring-integration-core/src/test/java/org/springframework/integration/graph/IntegrationGraphServerTests.java index 1cc4462a03..4ff419222c 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/graph/IntegrationGraphServerTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/graph/IntegrationGraphServerTests.java @@ -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); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/message/MethodInvokingMessageHandlerTests.java b/spring-integration-core/src/test/java/org/springframework/integration/message/MethodInvokingMessageHandlerTests.java index 88453756fb..531f93695d 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/message/MethodInvokingMessageHandlerTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/message/MethodInvokingMessageHandlerTests.java @@ -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 queue; - - TestBean(BlockingQueue queue) { - this.queue = queue; - } + private record TestBean(BlockingQueue queue) { @SuppressWarnings("unused") public void foo(String s) { diff --git a/spring-integration-core/src/test/java/org/springframework/integration/transformer/ContentEnricherTests.java b/spring-integration-core/src/test/java/org/springframework/integration/transformer/ContentEnricherTests.java index cec4c66d04..dc2ea127f0 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/transformer/ContentEnricherTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/transformer/ContentEnricherTests.java @@ -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 expressions = new HashMap(); + final Map 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 propertyExpressions = new HashMap(); + Map 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 propertyExpressions = new HashMap(); + Map 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 propertyExpressions = new HashMap(); + Map 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 propertyExpressions = new HashMap(); + Map 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 propertyExpressions = new HashMap(); + Map 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 propertyExpressions = new HashMap(); + Map 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; diff --git a/spring-integration-ftp/src/test/java/org/springframework/integration/ftp/inbound/FtpStreamingMessageSourceTests.java b/spring-integration-ftp/src/test/java/org/springframework/integration/ftp/inbound/FtpStreamingMessageSourceTests.java index 2932770233..b50b9da40d 100644 --- a/spring-integration-ftp/src/test/java/org/springframework/integration/ftp/inbound/FtpStreamingMessageSourceTests.java +++ b/spring-integration-ftp/src/test/java/org/springframework/integration/ftp/inbound/FtpStreamingMessageSourceTests.java @@ -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; } diff --git a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/TcpInboundGateway.java b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/TcpInboundGateway.java index 9fd5f284a7..5eefd506ba 100644 --- a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/TcpInboundGateway.java +++ b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/TcpInboundGateway.java @@ -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)); } } diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNioConnectionReadTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNioConnectionReadTests.java index 0539f25961..7fc9613c35 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNioConnectionReadTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNioConnectionReadTests.java @@ -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> responses = new ArrayList>(); + final List> 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> responses = new ArrayList>(); + final List> 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> responses = new ArrayList>(); + final List> 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> responses = new ArrayList>(); + final List> 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 added = new ArrayList(); - final List removed = new ArrayList(); + final List added = new ArrayList<>(); + final List removed = new ArrayList<>(); final CountDownLatch errorMessageLetch = new CountDownLatch(1); - final AtomicReference errorMessageRef = new AtomicReference(); + final AtomicReference 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 added = new ArrayList(); - final List removed = new ArrayList(); + final List added = new ArrayList<>(); + final List removed = new ArrayList<>(); final CountDownLatch errorMessageLetch = new CountDownLatch(1); - final AtomicReference errorMessageRef = new AtomicReference(); + final AtomicReference 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 added = new ArrayList(); - final List removed = new ArrayList(); + final List added = new ArrayList<>(); + final List removed = new ArrayList<>(); final CountDownLatch errorMessageLetch = new CountDownLatch(1); - final AtomicReference errorMessageRef = new AtomicReference(); + final AtomicReference 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 added = new ArrayList(); - final List removed = new ArrayList(); + final List added = new ArrayList<>(); + final List removed = new ArrayList<>(); final CountDownLatch errorMessageLetch = new CountDownLatch(1); - final AtomicReference errorMessageRef = new AtomicReference(); + final AtomicReference 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 added = new ArrayList(); - final List removed = new ArrayList(); + final List added = new ArrayList<>(); + final List removed = new ArrayList<>(); final CountDownLatch errorMessageLetch = new CountDownLatch(1); - final AtomicReference errorMessageRef = new AtomicReference(); + final AtomicReference 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); diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/util/SocketTestUtils.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/util/SocketTestUtils.java index 9dd796630f..0539b4018c 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/util/SocketTestUtils.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/util/SocketTestUtils.java @@ -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); diff --git a/spring-integration-mail/src/main/java/org/springframework/integration/mail/ImapIdleChannelAdapter.java b/spring-integration-mail/src/main/java/org/springframework/integration/mail/ImapIdleChannelAdapter.java index 59656b1831..238dafd649 100755 --- a/spring-integration-mail/src/main/java/org/springframework/integration/mail/ImapIdleChannelAdapter.java +++ b/spring-integration-mail/src/main/java/org/springframework/integration/mail/ImapIdleChannelAdapter.java @@ -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); } } diff --git a/spring-integration-mail/src/main/java/org/springframework/integration/mail/ImapMailReceiver.java b/spring-integration-mail/src/main/java/org/springframework/integration/mail/ImapMailReceiver.java index 293f0d2a43..acd2cab6e2 100755 --- a/spring-integration-mail/src/main/java/org/springframework/integration/mail/ImapMailReceiver.java +++ b/spring-integration-mail/src/main/java/org/springframework/integration/mail/ImapMailReceiver.java @@ -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); } diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/MqttPahoMessageDrivenChannelAdapter.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/MqttPahoMessageDrivenChannelAdapter.java index 5978d74c2c..19c88b6da6 100644 --- a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/MqttPahoMessageDrivenChannelAdapter.java +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/MqttPahoMessageDrivenChannelAdapter.java @@ -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"); diff --git a/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/MqttAdapterTests.java b/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/MqttAdapterTests.java index 9569fefc1d..a8c742017a 100644 --- a/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/MqttAdapterTests.java +++ b/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/MqttAdapterTests.java @@ -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 diff --git a/spring-integration-sftp/src/test/java/org/springframework/integration/sftp/inbound/SftpStreamingMessageSourceTests.java b/spring-integration-sftp/src/test/java/org/springframework/integration/sftp/inbound/SftpStreamingMessageSourceTests.java index 7425d2814c..8eebc73383 100644 --- a/spring-integration-sftp/src/test/java/org/springframework/integration/sftp/inbound/SftpStreamingMessageSourceTests.java +++ b/spring-integration-sftp/src/test/java/org/springframework/integration/sftp/inbound/SftpStreamingMessageSourceTests.java @@ -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; } diff --git a/spring-integration-stomp/src/main/java/org/springframework/integration/stomp/AbstractStompSessionManager.java b/spring-integration-stomp/src/main/java/org/springframework/integration/stomp/AbstractStompSessionManager.java index 04d6c79297..b6d66ac3fa 100644 --- a/spring-integration-stomp/src/main/java/org/springframework/integration/stomp/AbstractStompSessionManager.java +++ b/spring-integration-stomp/src/main/java/org/springframework/integration/stomp/AbstractStompSessionManager.java @@ -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() { + new ListenableFutureCallback<>() { @Override public void onFailure(Throwable ex) { diff --git a/spring-integration-stream/src/test/java/org/springframework/integration/stream/CharacterStreamSourceTests.java b/spring-integration-stream/src/test/java/org/springframework/integration/stream/CharacterStreamSourceTests.java index cd92532e38..8cf678ec47 100644 --- a/spring-integration-stream/src/test/java/org/springframework/integration/stream/CharacterStreamSourceTests.java +++ b/spring-integration-stream/src/test/java/org/springframework/integration/stream/CharacterStreamSourceTests.java @@ -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)); diff --git a/spring-integration-stream/src/test/java/org/springframework/integration/stream/CharacterStreamWritingMessageHandlerTests.java b/spring-integration-stream/src/test/java/org/springframework/integration/stream/CharacterStreamWritingMessageHandlerTests.java index 7e3737ef27..58cd098c16 100644 --- a/spring-integration-stream/src/test/java/org/springframework/integration/stream/CharacterStreamWritingMessageHandlerTests.java +++ b/spring-integration-stream/src/test/java/org/springframework/integration/stream/CharacterStreamWritingMessageHandlerTests.java @@ -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("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("foo"), 0); - channel.send(new GenericMessage("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("foo"), 0); - channel.send(new GenericMessage("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("foo"), 0); - channel.send(new GenericMessage("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("foo"), 0); - channel.send(new GenericMessage("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)); + 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(testObject1), 0); - channel.send(new GenericMessage(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(testObject1), 0); - channel.send(new GenericMessage(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; diff --git a/spring-integration-test-support/src/main/java/org/springframework/integration/test/util/OnlyOnceTrigger.java b/spring-integration-test-support/src/main/java/org/springframework/integration/test/util/OnlyOnceTrigger.java index 06bbe33deb..753588cd63 100644 --- a/spring-integration-test-support/src/main/java/org/springframework/integration/test/util/OnlyOnceTrigger.java +++ b/spring-integration-test-support/src/main/java/org/springframework/integration/test/util/OnlyOnceTrigger.java @@ -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); } } diff --git a/spring-integration-ws/src/test/java/org/springframework/integration/ws/config/WebServiceOutboundGatewayParserTests.java b/spring-integration-ws/src/test/java/org/springframework/integration/ws/config/WebServiceOutboundGatewayParserTests.java index 33e5b9b192..585139b552 100644 --- a/spring-integration-ws/src/test/java/org/springframework/integration/ws/config/WebServiceOutboundGatewayParserTests.java +++ b/spring-integration-ws/src/test/java/org/springframework/integration/ws/config/WebServiceOutboundGatewayParserTests.java @@ -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 diff --git a/src/reference/asciidoc/amqp.adoc b/src/reference/asciidoc/amqp.adoc index f3e6145c19..ecde6c4b24 100644 --- a/src/reference/asciidoc/amqp.adoc +++ b/src/reference/asciidoc/amqp.adoc @@ -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. diff --git a/src/reference/asciidoc/endpoint.adoc b/src/reference/asciidoc/endpoint.adoc index b8b9f71b6c..b9e6ed0dcd 100644 --- a/src/reference/asciidoc/endpoint.adoc +++ b/src/reference/asciidoc/endpoint.adoc @@ -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); ---- ====