From d766d233f5d15734521afbcba9dcfdf6e207c117 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 21 Jun 2017 12:12:49 -0400 Subject: [PATCH] Fix ReactiveStreamsConsumerTests and Checkstyle Also increase timeouts in the `RedisAvailableRule` --- .../channel/FluxMessageChannel.java | 3 +- .../ReactiveStreamsConsumerTests.java | 30 ++++++++++--------- .../redis/rules/RedisAvailableRule.java | 4 +-- 3 files changed, 20 insertions(+), 17 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/FluxMessageChannel.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/FluxMessageChannel.java index b3ff9104e3..6b63d2c600 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/channel/FluxMessageChannel.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/FluxMessageChannel.java @@ -60,7 +60,8 @@ public class FluxMessageChannel extends AbstractMessageChannel @Override protected boolean doSend(Message message, long timeout) { - Assert.state(subscribers.size() > 0, () -> "The [" + this + "] doesn't have subscribers to accept messages"); + Assert.state(this.subscribers.size() > 0, + () -> "The [" + this + "] doesn't have subscribers to accept messages"); this.sink.next(message); return true; } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/ReactiveStreamsConsumerTests.java b/spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/ReactiveStreamsConsumerTests.java index c1ab7254ef..bc51a5d90c 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/ReactiveStreamsConsumerTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/ReactiveStreamsConsumerTests.java @@ -88,7 +88,14 @@ public class ReactiveStreamsConsumerTests { reactiveConsumer.stop(); - testChannel.send(testMessage); + try { + testChannel.send(testMessage); + } + catch (Exception e) { + assertThat(e, instanceOf(MessageDeliveryException.class)); + assertThat(e.getCause(), instanceOf(IllegalStateException.class)); + assertThat(e.getMessage(), containsString("doesn't have subscribers to accept messages")); + } reactiveConsumer.start(); @@ -246,7 +253,14 @@ public class ReactiveStreamsConsumerTests { endpointFactoryBean.stop(); - testChannel.send(testMessage); + try { + testChannel.send(testMessage); + } + catch (Exception e) { + assertThat(e, instanceOf(MessageDeliveryException.class)); + assertThat(e.getCause(), instanceOf(IllegalStateException.class)); + assertThat(e.getMessage(), containsString("doesn't have subscribers to accept messages")); + } endpointFactoryBean.start(); @@ -260,16 +274,4 @@ public class ReactiveStreamsConsumerTests { assertThat(result, Matchers.>contains(testMessage, testMessage2, testMessage2)); } - @Test - public void testFluxMessageChannelSendWithoutSubscription() { - try { - new FluxMessageChannel().send(new GenericMessage<>("foo")); - } - catch (Exception e) { - assertThat(e, instanceOf(MessageDeliveryException.class)); - assertThat(e.getCause(), instanceOf(IllegalStateException.class)); - assertThat(e.getMessage(), containsString("doesn't have subscribers to accept messages")); - } - } - } diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/rules/RedisAvailableRule.java b/spring-integration-redis/src/test/java/org/springframework/integration/redis/rules/RedisAvailableRule.java index 77c44f5754..474e9b870b 100644 --- a/spring-integration-redis/src/test/java/org/springframework/integration/redis/rules/RedisAvailableRule.java +++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/rules/RedisAvailableRule.java @@ -47,8 +47,8 @@ public final class RedisAvailableRule implements MethodRule { redisStandaloneConfiguration.setPort(REDIS_PORT); JedisClientConfiguration clientConfiguration = JedisClientConfiguration.builder() - .connectTimeout(Duration.ofSeconds(10)) - .readTimeout(Duration.ofSeconds(10)) + .connectTimeout(Duration.ofSeconds(20)) + .readTimeout(Duration.ofSeconds(20)) .build(); connectionFactory = new JedisConnectionFactory(redisStandaloneConfiguration, clientConfiguration);