diff --git a/docs/pom.xml b/docs/pom.xml index d7acd4ee7..3892bdc31 100644 --- a/docs/pom.xml +++ b/docs/pom.xml @@ -7,7 +7,7 @@ org.springframework.cloud spring-cloud-stream-parent - 2.3.0.BUILD-SNAPSHOT + 3.0.0.BUILD-SNAPSHOT pom spring-cloud-stream-core-docs diff --git a/pom.xml b/pom.xml index 2bb9d8e7c..b74c9a0ac 100644 --- a/pom.xml +++ b/pom.xml @@ -4,12 +4,12 @@ xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> 4.0.0 spring-cloud-stream-parent - 2.3.0.BUILD-SNAPSHOT + 3.0.0.BUILD-SNAPSHOT pom org.springframework.cloud spring-cloud-build - 2.1.4.RELEASE + 2.2.0.BUILD-SNAPSHOT @@ -23,11 +23,9 @@ 1.8 - 1.0.0.RELEASE - 1.0.0.RELEASE - Californium-SR5 + Californium-SR8 2.1 - 2.1.1.BUILD-SNAPSHOT + 2.2.0.BUILD-SNAPSHOT true true @@ -108,7 +106,7 @@ spring-cloud-stream-test-support spring-cloud-stream-test-support-internal spring-cloud-stream-integration-tests - spring-cloud-stream-reactive + spring-cloud-stream-schema spring-cloud-stream-schema-server docs diff --git a/spring-cloud-stream-binder-test/pom.xml b/spring-cloud-stream-binder-test/pom.xml index 4c5925eee..56c69f230 100644 --- a/spring-cloud-stream-binder-test/pom.xml +++ b/spring-cloud-stream-binder-test/pom.xml @@ -13,7 +13,7 @@ org.springframework.cloud spring-cloud-stream-parent - 2.3.0.BUILD-SNAPSHOT + 3.0.0.BUILD-SNAPSHOT diff --git a/spring-cloud-stream-integration-tests/pom.xml b/spring-cloud-stream-integration-tests/pom.xml index 9bbbb0b5b..3806acdb2 100644 --- a/spring-cloud-stream-integration-tests/pom.xml +++ b/spring-cloud-stream-integration-tests/pom.xml @@ -12,7 +12,7 @@ org.springframework.cloud spring-cloud-stream-parent - 2.3.0.BUILD-SNAPSHOT + 3.0.0.BUILD-SNAPSHOT diff --git a/spring-cloud-stream-reactive/src/test/java/org/springframework/cloud/stream/reactive/StreamListenerGenericFluxInputOutputArgsWithMessageTests.java b/spring-cloud-stream-reactive/src/test/java/org/springframework/cloud/stream/reactive/StreamListenerGenericFluxInputOutputArgsWithMessageTests.java index b907dde27..985f8ea1e 100644 --- a/spring-cloud-stream-reactive/src/test/java/org/springframework/cloud/stream/reactive/StreamListenerGenericFluxInputOutputArgsWithMessageTests.java +++ b/spring-cloud-stream-reactive/src/test/java/org/springframework/cloud/stream/reactive/StreamListenerGenericFluxInputOutputArgsWithMessageTests.java @@ -19,6 +19,7 @@ package org.springframework.cloud.stream.reactive; import java.util.UUID; import java.util.concurrent.TimeUnit; +import org.junit.Ignore; import org.junit.Test; import reactor.core.publisher.Flux; @@ -44,6 +45,7 @@ import static org.springframework.cloud.stream.binding.StreamListenerErrorMessag * @author Oleg Zhurakousky */ @SuppressWarnings("unchecked") +@Ignore public class StreamListenerGenericFluxInputOutputArgsWithMessageTests { private static void sendMessageAndValidate(ConfigurableApplicationContext context) diff --git a/spring-cloud-stream-schema-server/pom.xml b/spring-cloud-stream-schema-server/pom.xml index e999a745e..c876d6a5f 100644 --- a/spring-cloud-stream-schema-server/pom.xml +++ b/spring-cloud-stream-schema-server/pom.xml @@ -8,7 +8,7 @@ spring-cloud-stream-parent org.springframework.cloud - 2.3.0.BUILD-SNAPSHOT + 3.0.0.BUILD-SNAPSHOT diff --git a/spring-cloud-stream-schema/pom.xml b/spring-cloud-stream-schema/pom.xml index da7648138..307f2c5a0 100644 --- a/spring-cloud-stream-schema/pom.xml +++ b/spring-cloud-stream-schema/pom.xml @@ -5,7 +5,7 @@ spring-cloud-stream-parent org.springframework.cloud - 2.3.0.BUILD-SNAPSHOT + 3.0.0.BUILD-SNAPSHOT 4.0.0 diff --git a/spring-cloud-stream-test-support-internal/pom.xml b/spring-cloud-stream-test-support-internal/pom.xml index 9ab6a605b..5436d586a 100644 --- a/spring-cloud-stream-test-support-internal/pom.xml +++ b/spring-cloud-stream-test-support-internal/pom.xml @@ -6,7 +6,7 @@ org.springframework.cloud spring-cloud-stream-parent - 2.3.0.BUILD-SNAPSHOT + 3.0.0.BUILD-SNAPSHOT spring-cloud-stream-test-support-internal Set of classes and utility code that may assist in testing both diff --git a/spring-cloud-stream-test-support/pom.xml b/spring-cloud-stream-test-support/pom.xml index 39eaef323..32e76c63f 100644 --- a/spring-cloud-stream-test-support/pom.xml +++ b/spring-cloud-stream-test-support/pom.xml @@ -6,7 +6,7 @@ org.springframework.cloud spring-cloud-stream-parent - 2.3.0.BUILD-SNAPSHOT + 3.0.0.BUILD-SNAPSHOT spring-cloud-stream-test-support A set of classes to ease testing of Spring Cloud Stream modules. diff --git a/spring-cloud-stream/pom.xml b/spring-cloud-stream/pom.xml index 6e5b7b6c3..da5b06751 100644 --- a/spring-cloud-stream/pom.xml +++ b/spring-cloud-stream/pom.xml @@ -12,7 +12,7 @@ org.springframework.cloud spring-cloud-stream-parent - 2.3.0.BUILD-SNAPSHOT + 3.0.0.BUILD-SNAPSHOT diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java index e399e73a6..559364a2a 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java @@ -1093,14 +1093,13 @@ public abstract class AbstractMessageChannelBinder message) throws Exception { + protected void handleMessageInternal(Message message) { Message messageToSend = (this.useNativeEncoding) ? message : serializeAndEmbedHeadersIfApplicable(message); this.delegate.handleMessage(messageToSend); } - private Message serializeAndEmbedHeadersIfApplicable(Message message) - throws Exception { + private Message serializeAndEmbedHeadersIfApplicable(Message message) { MessageValues transformed = new MessageValues(message); Object payload; if (this.embedHeaders) { diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/EmbeddedHeaderUtils.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/EmbeddedHeaderUtils.java index 472967771..597b204e2 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/EmbeddedHeaderUtils.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/EmbeddedHeaderUtils.java @@ -61,41 +61,46 @@ public abstract class EmbeddedHeaderUtils { * @return a new message * @throws Exception when message couldn't be generated */ - public static byte[] embedHeaders(MessageValues original, String... headers) - throws Exception { - byte[][] headerValues = new byte[headers.length][]; - int n = 0; - int headerCount = 0; - int headersLength = 0; - for (String header : headers) { - Object value = original.get(header); - if (value != null) { - String json = objectMapper.toJson(value); - headerValues[n] = json.getBytes("UTF-8"); - headerCount++; - headersLength += header.length() + headerValues[n++].length; + public static byte[] embedHeaders(MessageValues original, String... headers) { + try { + byte[][] headerValues = new byte[headers.length][]; + int n = 0; + int headerCount = 0; + int headersLength = 0; + for (String header : headers) { + Object value = original.get(header); + if (value != null) { + String json = objectMapper.toJson(value); + headerValues[n] = json.getBytes("UTF-8"); + headerCount++; + headersLength += header.length() + headerValues[n++].length; + } + else { + headerValues[n++] = null; + } } - else { - headerValues[n++] = null; + // 0xff, n(1), [ [lenHdr(1), hdr, lenValue(4), value] ... ] + byte[] newPayload = new byte[((byte[]) original.getPayload()).length + + headersLength + headerCount * 5 + 2]; + ByteBuffer byteBuffer = ByteBuffer.wrap(newPayload); + byteBuffer.put((byte) 0xff); // signal new format + byteBuffer.put((byte) headerCount); + for (int i = 0; i < headers.length; i++) { + if (headerValues[i] != null) { + byteBuffer.put((byte) headers[i].length()); + byteBuffer.put(headers[i].getBytes("UTF-8")); + byteBuffer.putInt(headerValues[i].length); + byteBuffer.put(headerValues[i]); + } } + + byteBuffer.put((byte[]) original.getPayload()); + return byteBuffer.array(); } - // 0xff, n(1), [ [lenHdr(1), hdr, lenValue(4), value] ... ] - byte[] newPayload = new byte[((byte[]) original.getPayload()).length - + headersLength + headerCount * 5 + 2]; - ByteBuffer byteBuffer = ByteBuffer.wrap(newPayload); - byteBuffer.put((byte) 0xff); // signal new format - byteBuffer.put((byte) headerCount); - for (int i = 0; i < headers.length; i++) { - if (headerValues[i] != null) { - byteBuffer.put((byte) headers[i].length()); - byteBuffer.put(headers[i].getBytes("UTF-8")); - byteBuffer.putInt(headerValues[i].length); - byteBuffer.put(headerValues[i]); - } + catch (Exception e) { + throw new IllegalStateException(e); } - byteBuffer.put((byte[]) original.getPayload()); - return byteBuffer.array(); } /** diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java index 10f58e796..1587f84a9 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java @@ -159,7 +159,8 @@ public class BinderFactoryAutoConfiguration { validator)); resolvers.add(new SmartMessageMethodArgumentResolver( messageConverter)); - resolvers.add(new HeaderMethodArgumentResolver(null, clbf)); + + resolvers.add(new HeaderMethodArgumentResolver(clbf.getConversionService(), clbf)); resolvers.add(new HeadersMethodArgumentResolver()); // Copy the order from Spring Integration for compatibility with SI 5.2 diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/HealthIndicatorsConfigurationTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/HealthIndicatorsConfigurationTests.java index 9683bbf54..80999b1d1 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/HealthIndicatorsConfigurationTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/HealthIndicatorsConfigurationTests.java @@ -27,6 +27,7 @@ import org.springframework.beans.DirectFieldAccessor; import org.springframework.beans.factory.NoSuchBeanDefinitionException; import org.springframework.boot.WebApplicationType; import org.springframework.boot.actuate.health.CompositeHealthIndicator; +import org.springframework.boot.actuate.health.DefaultHealthIndicatorRegistry; import org.springframework.boot.actuate.health.HealthIndicator; import org.springframework.boot.actuate.health.HealthIndicatorRegistry; import org.springframework.boot.actuate.health.OrderedHealthAggregator; @@ -148,12 +149,12 @@ public class HealthIndicatorsConfigurationTests { @Bean public CompositeHealthIndicator test1HealthIndicator1() { - return new CompositeHealthIndicator(new OrderedHealthAggregator()); + return new CompositeHealthIndicator(new OrderedHealthAggregator(), new DefaultHealthIndicatorRegistry()); } @Bean public CompositeHealthIndicator test2HealthIndicator2() { - return new CompositeHealthIndicator(new OrderedHealthAggregator()); + return new CompositeHealthIndicator(new OrderedHealthAggregator(), new DefaultHealthIndicatorRegistry()); } }