From ea5bbe83e5ac0d81701c31de4e3e38c12a50c4d3 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Thu, 18 Aug 2016 15:33:24 -0400 Subject: [PATCH] INT-4062: Reject Exec. Channel Early Subscription JIRA: https://jira.spring.io/browse/INT-4062 Only Enforce When an Executor is Provided With `ExecutorChannel` and `PublishSubscribeChannel` (when an executor is provided), any early subscription is lost because the dispatcher is replaced. checkstyle --- .../integration/channel/ExecutorChannel.java | 11 ++-- .../channel/PublishSubscribeChannel.java | 5 +- .../channel/ExecutorChannelTests.java | 18 +++++++ .../channel/PublishSubscribeChannelTests.java | 52 +++++++++++++++++++ 4 files changed, 78 insertions(+), 8 deletions(-) create mode 100644 spring-integration-core/src/test/java/org/springframework/integration/channel/PublishSubscribeChannelTests.java diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/ExecutorChannel.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/ExecutorChannel.java index a3730f914c..e780579d33 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/channel/ExecutorChannel.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/ExecutorChannel.java @@ -100,13 +100,10 @@ public class ExecutorChannel extends AbstractExecutorChannel { } @Override - public final void onInit() { - try { - super.onInit(); // TODO add throws clause in 5.0 - } - catch (Exception e) { - throw new IllegalStateException(e); - } + public final void onInit() throws Exception { + Assert.state(getDispatcher().getHandlerCount() == 0, "You cannot subscribe() until the channel " + + "bean is fully initialized by the framework. Do not subscribe in a @Bean definition"); + super.onInit(); if (!(this.executor instanceof ErrorHandlingTaskExecutor)) { ErrorHandler errorHandler = new MessagePublishingErrorHandler( new BeanFactoryChannelResolver(this.getBeanFactory())); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/PublishSubscribeChannel.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/PublishSubscribeChannel.java index 86583239bd..f377e920c7 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/channel/PublishSubscribeChannel.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/PublishSubscribeChannel.java @@ -24,6 +24,7 @@ import org.springframework.integration.dispatcher.MessageHandlingTaskDecorator; import org.springframework.integration.support.channel.BeanFactoryChannelResolver; import org.springframework.integration.util.ErrorHandlingTaskExecutor; import org.springframework.messaging.support.MessageHandlingRunnable; +import org.springframework.util.Assert; import org.springframework.util.ErrorHandler; /** @@ -135,6 +136,9 @@ public class PublishSubscribeChannel extends AbstractExecutorChannel { public final void onInit() throws Exception { super.onInit(); if (this.executor != null) { + Assert.state(getDispatcher().getHandlerCount() == 0, + "When providing an Executor, you cannot subscribe() until the channel " + + "bean is fully initialized by the framework. Do not subscribe in a @Bean definition"); if (!(this.executor instanceof ErrorHandlingTaskExecutor)) { if (this.errorHandler == null) { this.errorHandler = new MessagePublishingErrorHandler( @@ -167,7 +171,6 @@ public class PublishSubscribeChannel extends AbstractExecutorChannel { } }); - } @Override diff --git a/spring-integration-core/src/test/java/org/springframework/integration/channel/ExecutorChannelTests.java b/spring-integration-core/src/test/java/org/springframework/integration/channel/ExecutorChannelTests.java index 669bff6d68..66c8552439 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/channel/ExecutorChannelTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/channel/ExecutorChannelTests.java @@ -16,17 +16,21 @@ package org.springframework.integration.channel; +import static org.hamcrest.Matchers.equalTo; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertNull; import static org.junit.Assert.assertSame; +import static org.junit.Assert.assertThat; import static org.junit.Assert.assertTrue; +import static org.junit.Assert.fail; import static org.mockito.BDDMockito.willThrow; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.verify; import java.util.concurrent.CountDownLatch; +import java.util.concurrent.Executor; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; @@ -210,6 +214,20 @@ public class ExecutorChannelTests { assertTrue(interceptor.wasAfterHandledInvoked()); } + @Test + public void testEarlySubscribe() { + ExecutorChannel channel = new ExecutorChannel(mock(Executor.class)); + try { + channel.subscribe(m -> { }); + channel.setBeanFactory(mock(BeanFactory.class)); + channel.afterPropertiesSet(); + fail("expected Exception"); + } + catch (IllegalStateException e) { + assertThat(e.getMessage(), equalTo("You cannot subscribe() until the channel " + + "bean is fully initialized by the framework. Do not subscribe in a @Bean definition")); + } + } private static class TestHandler implements MessageHandler { diff --git a/spring-integration-core/src/test/java/org/springframework/integration/channel/PublishSubscribeChannelTests.java b/spring-integration-core/src/test/java/org/springframework/integration/channel/PublishSubscribeChannelTests.java new file mode 100644 index 0000000000..87b211a3c2 --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/channel/PublishSubscribeChannelTests.java @@ -0,0 +1,52 @@ +/* + * Copyright 2016 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. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.integration.channel; + +import static org.hamcrest.Matchers.equalTo; +import static org.junit.Assert.assertThat; +import static org.junit.Assert.fail; +import static org.mockito.Mockito.mock; + +import java.util.concurrent.Executor; + +import org.junit.Test; + +import org.springframework.beans.factory.BeanFactory; + +/** + * @author Gary Russell + * @since 5.0 + * + */ +public class PublishSubscribeChannelTests { + + @Test + public void testEarlySubscribe() { + PublishSubscribeChannel channel = new PublishSubscribeChannel(mock(Executor.class)); + try { + channel.subscribe(m -> { }); + channel.setBeanFactory(mock(BeanFactory.class)); + channel.afterPropertiesSet(); + fail("expected Exception"); + } + catch (IllegalStateException e) { + assertThat(e.getMessage(), equalTo("When providing an Executor, you cannot subscribe() until the channel " + + "bean is fully initialized by the framework. Do not subscribe in a @Bean definition")); + } + } + +}