Adding Spring Cloud Stream Version To Message Headers For Easier Debugging Of Issues.

Fix checkstyles
Resolves #3027
This commit is contained in:
Ömer Çelik
2024-11-04 13:23:01 +03:00
committed by Oleg Zhurakousky
parent 17ac6352e4
commit a8fb34db0e
7 changed files with 442 additions and 23 deletions

1
.gitignore vendored
View File

@@ -25,6 +25,7 @@ dump.rdb
.apt_generated
artifacts
**/dependency-reduced-pom.xml
core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/utils/GeneratedBuildProperties.java
node
node_modules

View File

@@ -16,7 +16,10 @@
package org.springframework.cloud.stream.function;
import java.nio.charset.StandardCharsets;
import java.util.Collections;
import java.util.Locale;
import java.util.Map;
import java.util.function.Function;
import org.junit.jupiter.api.BeforeAll;
@@ -26,9 +29,14 @@ import org.springframework.boot.WebApplicationType;
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
import org.springframework.boot.builder.SpringApplicationBuilder;
import org.springframework.cloud.function.context.message.MessageUtils;
import org.springframework.cloud.function.json.JsonMapper;
import org.springframework.cloud.stream.binder.BinderHeaders;
import org.springframework.cloud.stream.binder.test.EnableTestBinder;
import org.springframework.cloud.stream.binder.test.InputDestination;
import org.springframework.cloud.stream.binder.test.OutputDestination;
import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration;
import org.springframework.cloud.stream.utils.BuildInformationProvider;
import org.springframework.context.ApplicationContext;
import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@@ -43,7 +51,6 @@ import static org.assertj.core.api.Assertions.assertThat;
/**
* @author Omer Celik
*/
public class HeaderTests {
@BeforeAll
@@ -63,10 +70,8 @@ public class HeaderTests {
OutputDestination outputDestination = context.getBean(OutputDestination.class);
Message<byte[]> messageReceived = outputDestination.receive(1000, "emptyConfigurationDestination");
MessageHeaders headers = messageReceived.getHeaders();
assertThat(headers).isNotNull();
assertThat(headers.get(MessageUtils.TARGET_PROTOCOL)).isEqualTo("kafka");
assertThat(headers.get(MessageHeaders.CONTENT_TYPE)).isEqualTo("application/json");
checkCommonHeaders(messageReceived.getHeaders());
}
}
@@ -75,6 +80,7 @@ public class HeaderTests {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(EmptyConfiguration.class))
.web(WebApplicationType.NONE).run("--spring.jmx.enabled=false")) {
StreamBridge streamBridge = context.getBean(StreamBridge.class);
String jsonPayload = "{\"name\":\"Omer\"}";
streamBridge.send("myBinding-out-0",
@@ -82,13 +88,12 @@ public class HeaderTests {
.setHeader("anyHeader", "anyValue")
.build(),
MimeTypeUtils.APPLICATION_JSON);
OutputDestination output = context.getBean(OutputDestination.class);
Message<byte[]> result = output.receive(1000, "myBinding-out-0");
MessageHeaders headers = result.getHeaders();
assertThat(headers).isNotNull();
assertThat(headers.get(MessageUtils.TARGET_PROTOCOL)).isEqualTo("kafka");
assertThat(headers.get(MessageHeaders.CONTENT_TYPE)).isEqualTo("application/json");
assertThat(headers.get("anyHeader")).isEqualTo("anyValue");
checkCommonHeaders(result.getHeaders());
assertThat(result.getHeaders().get("anyHeader")).isEqualTo("anyValue");
}
}
@@ -99,16 +104,35 @@ public class HeaderTests {
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false",
"--spring.cloud.function.definition=uppercase")) {
String jsonPayload = "{\"surname\":\"Celik\"}";
InputDestination input = context.getBean(InputDestination.class);
input.send(new GenericMessage<>(jsonPayload.getBytes()), "uppercase-in-0");
OutputDestination output = context.getBean(OutputDestination.class);
OutputDestination output = context.getBean(OutputDestination.class);
Message<byte[]> result = output.receive(1000, "uppercase-out-0");
MessageHeaders headers = result.getHeaders();
assertThat(headers).isNotNull();
assertThat(headers.get(MessageUtils.TARGET_PROTOCOL)).isEqualTo("kafka");
assertThat(headers.get(MessageHeaders.CONTENT_TYPE)).isEqualTo("application/json");
checkCommonHeaders(result.getHeaders());
}
}
@Test
void checkGenericMessageSentUsingStreamBridge() {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(FunctionUpperCaseConfiguration.class))
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false",
"--spring.cloud.function.definition=uppercase")) {
String jsonPayload = "{\"anyFieldName\":\"anyValue\"}";
final StreamBridge streamBridge = context.getBean(StreamBridge.class);
GenericMessage<String> message = new GenericMessage<>(jsonPayload);
streamBridge.send("uppercase-in-0", message);
OutputDestination output = context.getBean(OutputDestination.class);
Message<byte[]> result = output.receive(1000, "uppercase-out-0");
checkCommonHeaders(result.getHeaders());
}
}
@@ -127,11 +151,96 @@ public class HeaderTests {
OutputDestination target = context.getBean(OutputDestination.class);
Message<byte[]> message = target.receive(5, "uppercase-out-0");
MessageHeaders headers = message.getHeaders();
checkCommonHeaders(message.getHeaders());
}
@Test
void checkStringToMapMessageStreamListener() {
ApplicationContext context = new SpringApplicationBuilder(
StringToMapMessageConfiguration.class).web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false");
InputDestination source = context.getBean(InputDestination.class);
String jsonPayload = "{\"name\":\"Omer\"}";
source.send(new GenericMessage<>(jsonPayload.getBytes()));
OutputDestination target = context.getBean(OutputDestination.class);
Message<byte[]> outputMessage = target.receive();
checkCommonHeaders(outputMessage.getHeaders());
}
@Test
void checkPojoToPojo() {
ApplicationContext context = new SpringApplicationBuilder(
PojoToPojoConfiguration.class).web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false");
InputDestination source = context.getBean(InputDestination.class);
String jsonPayload = "{\"name\":\"Omer\"}";
source.send(new GenericMessage<>(jsonPayload.getBytes()));
OutputDestination target = context.getBean(OutputDestination.class);
Message<byte[]> outputMessage = target.receive();
checkCommonHeaders(outputMessage.getHeaders());
}
@Test
void checkPojoToString() {
ApplicationContext context = new SpringApplicationBuilder(
PojoToStringConfiguration.class).web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false");
InputDestination source = context.getBean(InputDestination.class);
OutputDestination target = context.getBean(OutputDestination.class);
String jsonPayload = "{\"name\":\"Neso\"}";
source.send(new GenericMessage<>(jsonPayload.getBytes()));
Message<byte[]> outputMessage = target.receive();
checkCommonHeaders(outputMessage.getHeaders());
}
@Test
void checkPojoToByteArray() {
ApplicationContext context = new SpringApplicationBuilder(
PojoToByteArrayConfiguration.class).web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false");
InputDestination source = context.getBean(InputDestination.class);
OutputDestination target = context.getBean(OutputDestination.class);
String jsonPayload = "{\"name\":\"Neptune\"}";
source.send(new GenericMessage<>(jsonPayload.getBytes()));
Message<byte[]> outputMessage = target.receive();
checkCommonHeaders(outputMessage.getHeaders());
}
@Test
void checkStringToPojoInboundContentTypeHeader() {
ApplicationContext context = new SpringApplicationBuilder(
StringToPojoConfiguration.class).web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false");
InputDestination source = context.getBean(InputDestination.class);
OutputDestination target = context.getBean(OutputDestination.class);
String jsonPayload = "{\"name\":\"Mercury\"}";
source.send(new GenericMessage<>(jsonPayload.getBytes(),
new MessageHeaders(Collections.singletonMap(MessageHeaders.CONTENT_TYPE,
MimeTypeUtils.APPLICATION_JSON_VALUE))));
Message<byte[]> outputMessage = target.receive();
checkCommonHeaders(outputMessage.getHeaders());
}
@Test
void checkPojoMessageToStringMessage() {
ApplicationContext context = new SpringApplicationBuilder(
PojoMessageToStringMessageConfiguration.class)
.web(WebApplicationType.NONE).run("--spring.jmx.enabled=false");
InputDestination source = context.getBean(InputDestination.class);
OutputDestination target = context.getBean(OutputDestination.class);
String jsonPayload = "{\"name\":\"Earth\"}";
source.send(new GenericMessage<>(jsonPayload.getBytes()));
Message<byte[]> outputMessage = target.receive();
MessageHeaders headers = outputMessage.getHeaders();
assertThat(BuildInformationProvider.isVersionValid((String) headers.get(BinderHeaders.SCST_VERSION))).isTrue();
}
private void checkCommonHeaders(MessageHeaders headers) {
assertThat(headers).isNotNull();
assertThat(headers).isNotNull();
assertThat(headers.get(MessageHeaders.CONTENT_TYPE)).isEqualTo("application/json");
assertThat(headers.get(MessageUtils.TARGET_PROTOCOL)).isEqualTo("kafka");
assertThat(headers.get(MessageHeaders.CONTENT_TYPE)).isEqualTo("application/json");
assertThat(BuildInformationProvider.isVersionValid((String) headers.get(BinderHeaders.SCST_VERSION))).isTrue();
}
@EnableAutoConfiguration
@@ -156,6 +265,97 @@ public class HeaderTests {
}
}
@EnableTestBinder
@EnableAutoConfiguration
public static class StringToMapMessageConfiguration {
@Bean
public Function<Message<Map<?, ?>>, String> echo() {
return value -> {
assertThat(value.getPayload() instanceof Map).isTrue();
return (String) value.getPayload().get("name");
};
}
}
@EnableTestBinder
@EnableAutoConfiguration
public static class PojoToPojoConfiguration {
@Bean
public Function<Planet, Planet> echo() {
return value -> value;
}
}
@EnableTestBinder
@EnableAutoConfiguration
public static class PojoToStringConfiguration {
@Bean
public Function<Planet, String> echo() {
return Planet::toString;
}
}
@EnableTestBinder
@EnableAutoConfiguration
public static class PojoToByteArrayConfiguration {
@Bean
public Function<Planet, byte[]> echo() {
return value -> value.toString().getBytes(StandardCharsets.UTF_8);
}
}
@EnableTestBinder
@EnableAutoConfiguration
public static class StringToPojoConfiguration {
@Bean
public Function<String, Planet> echo(JsonMapper mapper) {
return value -> mapper.fromJson(value, Planet.class);
}
}
@EnableTestBinder
@EnableAutoConfiguration
public static class PojoMessageToStringMessageConfiguration {
@Bean
public Function<Message<Planet>, Message<String>> echo() {
return value -> MessageBuilder.withPayload(value.getPayload().toString())
.setHeader("expected-content-type", MimeTypeUtils.TEXT_PLAIN_VALUE)
.build();
}
}
public static class Planet {
private String name;
Planet() {
this(null);
}
Planet(String name) {
this.name = name;
}
public String getName() {
return this.name;
}
public void setName(String name) {
this.name = name;
}
@Override
public String toString() {
return this.name;
}
}
public static class EmptyPojo {
}

View File

@@ -15,6 +15,10 @@
<version>4.2.0-SNAPSHOT</version>
</parent>
<properties>
<timestamp>${maven.build.timestamp}</timestamp>
</properties>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
@@ -129,6 +133,77 @@
<jvmTarget>1.8</jvmTarget>
</configuration>
</plugin>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-clean-plugin</artifactId>
<executions>
<execution>
<phase>clean</phase>
<goals>
<goal>clean</goal>
</goals>
<configuration>
<filesets>
<fileset>
<directory>src/main/java</directory>
<includes>
<include>**/GeneratedBuildProperties.java</include>
</includes>
</fileset>
</filesets>
</configuration>
</execution>
</executions>
</plugin>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-antrun-plugin</artifactId>
<executions>
<execution>
<id>generate-build-info</id>
<phase>generate-sources</phase>
<goals>
<goal>run</goal>
</goals>
<configuration>
<target>
<copy file="src/main/template/org/springframework/cloud/stream/utils/GeneratedBuildProperties.java"
tofile="src/main/java/org/springframework/cloud/stream/utils/GeneratedBuildProperties.java">
<filterchain>
<replacestring from="@project.version@" to="${project.version}"/>
<replacestring from="@timestamp@" to="${timestamp}"/>
</filterchain>
</copy>
</target>
</configuration>
</execution>
</executions>
</plugin>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-resources-plugin</artifactId>
<executions>
<execution>
<phase>process-sources</phase>
<goals>
<goal>copy-resources</goal>
</goals>
<configuration>
<resources>
<resource>
<directory>src/main/template</directory>
<includes>
<include>**/*.java</include>
</includes>
<filtering>true</filtering>
</resource>
</resources>
<outputDirectory>src/main/java</outputDirectory>
<overwrite>true</overwrite>
</configuration>
</execution>
</executions>
</plugin>
</plugins>
</build>

View File

@@ -79,6 +79,7 @@ import org.springframework.cloud.stream.config.BindingProperties;
import org.springframework.cloud.stream.config.BindingServiceConfiguration;
import org.springframework.cloud.stream.config.BindingServiceProperties;
import org.springframework.cloud.stream.messaging.DirectWithAttributesChannel;
import org.springframework.cloud.stream.utils.BuildInformationProvider;
import org.springframework.context.ApplicationContext;
import org.springframework.context.ApplicationContextAware;
import org.springframework.context.ConfigurableApplicationContext;
@@ -132,6 +133,7 @@ import org.springframework.util.StringUtils;
* @author Ivan Shapoval
* @author Patrik Péter Süli
* @author Artem Bilan
* @author Omer Celik
* @since 2.1
*/
@Lazy(false)
@@ -470,7 +472,7 @@ public class FunctionConfiguration {
if (this.functionProperties.isComposeFrom()) {
AbstractSubscribableChannel outputChannel = this.applicationContext.getBean(outputBindingNames.iterator().next(), AbstractSubscribableChannel.class);
logger.info("Composing at the head of output destination: " + outputChannel.getBeanName());
String outputChannelName = ((AbstractMessageChannel) outputChannel).getBeanName();
String outputChannelName = outputChannel.getBeanName();
DirectWithAttributesChannel newOutputChannel = new DirectWithAttributesChannel();
newOutputChannel.setAttribute("type", "output");
newOutputChannel.setComponentName("output.extended");
@@ -497,11 +499,14 @@ public class FunctionConfiguration {
headersField.setAccessible(true);
targetProtocolEnhancer.set(message -> {
Map<String, Object> headersMap = (Map<String, Object>) ReflectionUtils
.getField(headersField, ((Message) message).getHeaders());
.getField(headersField, message.getHeaders());
headersMap.putIfAbsent(MessageUtils.TARGET_PROTOCOL, targetProtocol);
if (CloudEventMessageUtils.isCloudEvent((message))) {
headersMap.putIfAbsent(MessageUtils.MESSAGE_TYPE, CloudEventMessageUtils.CLOUDEVENT_VALUE);
}
if (BuildInformationProvider.isVersionValid()) {
headersMap.putIfAbsent(BinderHeaders.SCST_VERSION, BuildInformationProvider.getVersion());
}
return message;
});
}
@@ -836,6 +841,9 @@ public class FunctionConfiguration {
if (CloudEventMessageUtils.isCloudEvent(message)) {
headersMap.putIfAbsent(MessageUtils.MESSAGE_TYPE, CloudEventMessageUtils.CLOUDEVENT_VALUE);
}
if (BuildInformationProvider.isVersionValid()) {
headersMap.putIfAbsent(BinderHeaders.SCST_VERSION, BuildInformationProvider.getVersion());
}
}
}

View File

@@ -17,7 +17,6 @@
package org.springframework.cloud.stream.function;
import java.lang.reflect.Type;
import java.util.Collections;
import java.util.HashMap;
import java.util.LinkedHashMap;
import java.util.Map;
@@ -43,6 +42,7 @@ import org.springframework.cloud.function.context.message.MessageUtils;
import org.springframework.cloud.function.core.FunctionInvocationHelper;
import org.springframework.cloud.stream.binder.Binder;
import org.springframework.cloud.stream.binder.BinderFactory;
import org.springframework.cloud.stream.binder.BinderHeaders;
import org.springframework.cloud.stream.binder.ProducerProperties;
import org.springframework.cloud.stream.binding.BindingService;
import org.springframework.cloud.stream.binding.DefaultPartitioningInterceptor;
@@ -50,6 +50,7 @@ import org.springframework.cloud.stream.binding.NewDestinationBindingCallback;
import org.springframework.cloud.stream.config.BindingProperties;
import org.springframework.cloud.stream.config.BindingServiceProperties;
import org.springframework.cloud.stream.messaging.DirectWithAttributesChannel;
import org.springframework.cloud.stream.utils.BuildInformationProvider;
import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.core.ResolvableType;
import org.springframework.integration.channel.AbstractMessageChannel;
@@ -85,6 +86,7 @@ import org.springframework.util.StringUtils;
* @author Soby Chacko
* @author Byungjun You
* @author Michał Rowicki
* @author Omer Celik
* @since 3.0.3
*
*/
@@ -206,8 +208,7 @@ public final class StreamBridge implements StreamOperations, SmartInitializingSi
this.applicationContext.getBean(BinderFactory.class));
Message<?> messageToSend = data instanceof Message messageData
? MessageBuilder.fromMessage(messageData).setHeaderIfAbsent(MessageUtils.TARGET_PROTOCOL, targetType).build()
: new GenericMessage<>(data, Collections.singletonMap(MessageUtils.TARGET_PROTOCOL, targetType));
? createMessageWithHeader(messageData, targetType) : createGenericMessageWithHeader(data, targetType);
Message<?> resultMessage;
lock.lock();
@@ -228,6 +229,24 @@ public final class StreamBridge implements StreamOperations, SmartInitializingSi
return messageChannel.send(resultMessage);
}
private Message<?> createMessageWithHeader(Message<?> messageData, String targetType) {
MessageBuilder messageBuilder = MessageBuilder.fromMessage(messageData)
.copyHeaders(createHeaders(targetType));
return messageBuilder.build();
}
private Message<?> createGenericMessageWithHeader(Object data, String targetType) {
return new GenericMessage<>(data, createHeaders(targetType));
}
private Map<String, Object> createHeaders(String targetType) {
Map<String, Object> headers = new HashMap<>();
headers.put(MessageUtils.TARGET_PROTOCOL, targetType);
if (BuildInformationProvider.isVersionValid()) {
headers.put(BinderHeaders.SCST_VERSION, BuildInformationProvider.getVersion());
}
return headers;
}
private int hashProducerProperties(ProducerProperties producerProperties, String outputContentType) {
int hash = outputContentType.hashCode()
+ Boolean.hashCode(producerProperties.isUseNativeEncoding())

View File

@@ -0,0 +1,66 @@
/*
* Copyright 2024-2024 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
*
* https://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.cloud.stream.utils;
/**
* Provides information about current Spring Cloud Stream build.
*
* @author Omer Celik
*/
public final class BuildInformationProvider {
private static final String UNKNOWN_SCST_VERSION = "-1";
private static final BuildInformation BUILD_INFO_CACHE;
static {
BUILD_INFO_CACHE = createBuildInformation();
}
private BuildInformationProvider() {
}
public static boolean isVersionValid() {
return !getVersion().equals(UNKNOWN_SCST_VERSION);
}
public static boolean isVersionValid(String version) {
return !version.equals(UNKNOWN_SCST_VERSION);
}
public static String getVersion() {
return BUILD_INFO_CACHE.version();
}
// If you have a compilation error at GeneratedBuildProperties then run 'mvn clean install'
// the GeneratedBuildProperties class is generated at a compile-time
private static BuildInformation createBuildInformation() {
return new BuildInformation(calculateVersion(), GeneratedBuildProperties.TIMESTAMP);
}
// If you have a compilation error at GeneratedBuildProperties then run 'mvn clean install'
// the GeneratedBuildProperties class is generated at a compile-time
private static String calculateVersion() {
String version = GeneratedBuildProperties.VERSION;
if (version.startsWith("@") && version.endsWith("@")) {
return UNKNOWN_SCST_VERSION;
}
return version;
}
private record BuildInformation(String version, String timestamp) {
}
}

View File

@@ -0,0 +1,50 @@
/*
* Copyright 2024-2024 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
*
* https://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.cloud.stream.utils;
import javax.annotation.processing.Generated;
/**
* Exposes the Spring Cloud Stream Build Properties.
* This class is generated in a build-time from a template stored at
* src/main/template/org/springframework/cloud/stream/utils/GeneratedBuildProperties.
*
* Do not edit by hand as the changes will be overwritten in the next build.
* We use this method to support all build types. (Fat jar, shaded jar, war, etc.)
*
* WARNING: DO NOT CHANGE FIELD VALUES IN THE TEMPLATE.(For example: @project.version@)
* The fields are injected using the @.....@ keywords with the Maven Antrun Plugin.
*
* @author Omer Celik
* @since 4.2.0
*/
@Generated("")
public final class GeneratedBuildProperties {
/**
* Indicates the Spring Cloud Stream version.
*/
public static final String VERSION = "@project.version@";
/**
* Indicates the build time of the project.
*/
public static final String TIMESTAMP = "@timestamp@";
private GeneratedBuildProperties() {
}
}