From 0bbabaf7054169aec99db14c176f6603304e0b82 Mon Sep 17 00:00:00 2001 From: Marius Bogoevici Date: Fri, 8 Apr 2016 08:46:38 -0400 Subject: [PATCH] Make AggregateApplication package local Fixes #474 - Now the only API for aggregating is AggregateApplicationBuilder - add namespace tests --- .../aggregate/AggregateApplication.java | 18 ++++----- .../aggregation/ModuleAggregationTest.java | 37 +++++++++++++++++-- 2 files changed, 43 insertions(+), 12 deletions(-) diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/aggregate/AggregateApplication.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/aggregate/AggregateApplication.java index faed0c19e..68af8bcd1 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/aggregate/AggregateApplication.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/aggregate/AggregateApplication.java @@ -35,7 +35,7 @@ import org.springframework.messaging.SubscribableChannel; * @author Ilayaperumal Gopinathan * @author Venil Noronha */ -public class AggregateApplication { +class AggregateApplication { private static final String SPRING_CLOUD_STREAM_INTERNAL_PREFIX = "spring.cloud.stream.internal"; @@ -55,14 +55,14 @@ public class AggregateApplication { * * @return the resulting parent context for the aggregate */ - public static ConfigurableApplicationContext run(Class[] apps, String[] parentArgs, String[][] appArgs) { + static ConfigurableApplicationContext run(Class[] apps, String[] parentArgs, String[][] appArgs) { ConfigurableApplicationContext parentContext = createParentContext(parentArgs != null ? parentArgs : new String[0]); runEmbedded(parentContext, apps, appArgs); return parentContext; } - public static ConfigurableApplicationContext run(Class... apps) { + static ConfigurableApplicationContext run(Class... apps) { return run(apps, null, null); } @@ -73,7 +73,7 @@ public class AggregateApplication { * @param apps a list of classes, representing root context definitions for apps * @param args arguments for the apps */ - public static void runEmbedded(ConfigurableApplicationContext parentContext, + static void runEmbedded(ConfigurableApplicationContext parentContext, Class[] apps, String[][] args) { SharedChannelRegistry bean = parentContext.getBean(SharedChannelRegistry.class); prepareSharedChannelRegistry(bean, apps); @@ -81,7 +81,7 @@ public class AggregateApplication { createChildContexts(parentContext, apps, args); } - protected static ConfigurableApplicationContext createParentContext(String[] args) { + static ConfigurableApplicationContext createParentContext(String[] args) { SpringApplicationBuilder aggregatorParentConfiguration = new SpringApplicationBuilder(); aggregatorParentConfiguration .sources(AggregatorParentConfiguration.class) @@ -92,7 +92,7 @@ public class AggregateApplication { return aggregatorParentConfiguration.run(args); } - private static void createChildContexts(ConfigurableApplicationContext parentContext, + static void createChildContexts(ConfigurableApplicationContext parentContext, Class[] apps, String args[][]) { for (int i = apps.length - 1; i >= 0; i--) { String appClassName = apps[i].getName(); @@ -101,7 +101,7 @@ public class AggregateApplication { } } - protected static String getNamespace(String appClassName, int index) { + static String getNamespace(String appClassName, int index) { return appClassName + "_" + index; } @@ -117,7 +117,7 @@ public class AggregateApplication { .parent(applicationContext); } - protected static void prepareSharedChannelRegistry(SharedChannelRegistry sharedChannelRegistry, Class[] apps) { + static void prepareSharedChannelRegistry(SharedChannelRegistry sharedChannelRegistry, Class[] apps) { LinkedHashMap, String> appsToRegister = new LinkedHashMap<>(); for (int i = apps.length - 1; i >= 0; i--) { String appClassName = apps[i].getName(); @@ -126,7 +126,7 @@ public class AggregateApplication { prepareSharedChannelRegistry(sharedChannelRegistry, appsToRegister); } - protected static void prepareSharedChannelRegistry(SharedChannelRegistry sharedChannelRegistry, + static void prepareSharedChannelRegistry(SharedChannelRegistry sharedChannelRegistry, LinkedHashMap, String> appsWithNamespace) { int i = 0; SubscribableChannel sharedChannel = null; diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/aggregation/ModuleAggregationTest.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/aggregation/ModuleAggregationTest.java index d33c8307b..e34f28c61 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/aggregation/ModuleAggregationTest.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/aggregation/ModuleAggregationTest.java @@ -16,12 +16,11 @@ package org.springframework.cloud.stream.aggregation; +import static org.hamcrest.CoreMatchers.instanceOf; import static org.hamcrest.collection.IsCollectionWithSize.hasSize; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertThat; -import org.junit.Test; - import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.cloud.stream.aggregate.AggregateApplicationBuilder; import org.springframework.cloud.stream.aggregate.SharedChannelRegistry; @@ -31,6 +30,9 @@ import org.springframework.cloud.stream.messaging.Processor; import org.springframework.cloud.stream.messaging.Source; import org.springframework.cloud.stream.utils.MockBinderRegistryConfiguration; import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.messaging.MessageChannel; + +import org.junit.Test; /** * @author Marius Bogoevici @@ -39,11 +41,40 @@ public class ModuleAggregationTest { @Test public void testModuleAggregation() { - ConfigurableApplicationContext aggregatedApplicationContext = new AggregateApplicationBuilder(MockBinderRegistryConfiguration.class).from(TestSource.class).to(TestProcessor.class).run(); + ConfigurableApplicationContext aggregatedApplicationContext = + new AggregateApplicationBuilder(MockBinderRegistryConfiguration.class) + .from(TestSource.class) + .to(TestProcessor.class) + .run(); SharedChannelRegistry sharedChannelRegistry = aggregatedApplicationContext.getBean(SharedChannelRegistry.class); BindableChannelFactory channelFactory = aggregatedApplicationContext.getBean(BindableChannelFactory.class); assertNotNull(channelFactory); assertThat(sharedChannelRegistry.getAll().keySet(), hasSize(2)); + aggregatedApplicationContext.close(); + } + + @Test + public void testNamespaces() { + ConfigurableApplicationContext aggregatedApplicationContext = + new AggregateApplicationBuilder(MockBinderRegistryConfiguration.class) + .from(TestSource.class) + .namespace("foo") + .to(TestProcessor.class) + .namespace("bar") + .run(); + SharedChannelRegistry sharedChannelRegistry + = aggregatedApplicationContext.getBean(SharedChannelRegistry.class); + BindableChannelFactory channelFactory + = aggregatedApplicationContext.getBean(BindableChannelFactory.class); + Object fooOutput = sharedChannelRegistry.get("foo.output"); + assertNotNull(fooOutput); + assertThat(fooOutput, instanceOf(MessageChannel.class)); + Object barInput = sharedChannelRegistry.get("bar.input"); + assertNotNull(barInput); + assertThat(barInput, instanceOf(MessageChannel.class)); + assertNotNull(channelFactory); + assertThat(sharedChannelRegistry.getAll().keySet(), hasSize(2)); + aggregatedApplicationContext.close(); }