From a834ce9ffba646813ea14a73ec59c355126a32c2 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 21 Jun 2017 11:10:41 -0400 Subject: [PATCH] FluxMessageChannel.send: no subscribers assertion --- .../integration/channel/FluxMessageChannel.java | 2 ++ .../reactive/ReactiveStreamsConsumerTests.java | 13 +++++++++++++ 2 files changed, 15 insertions(+) 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 6b42b60a19..b3ff9104e3 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 @@ -25,6 +25,7 @@ import org.reactivestreams.Publisher; import org.reactivestreams.Subscriber; import org.springframework.messaging.Message; +import org.springframework.util.Assert; import reactor.core.publisher.ConnectableFlux; import reactor.core.publisher.Flux; @@ -59,6 +60,7 @@ 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"); 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 c616338ee9..c1ab7254ef 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 @@ -16,6 +16,7 @@ package org.springframework.integration.channel.reactive; +import static org.hamcrest.Matchers.containsString; import static org.hamcrest.Matchers.equalTo; import static org.hamcrest.Matchers.instanceOf; import static org.junit.Assert.assertSame; @@ -259,4 +260,16 @@ 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")); + } + } + }