Add additional test and error handling

Handle when functionToInvoke returns null

add link to issue

Signed-off-by: Steven Gantz <steven.p.gantz@gmail.com>
This commit is contained in:
Steven Gantz
2025-02-10 20:41:41 -05:00
parent 7178dbfdac
commit 25427bc718
2 changed files with 26 additions and 3 deletions

View File

@@ -37,6 +37,7 @@ import java.util.function.Supplier;
import java.util.stream.Collectors;
import java.util.stream.IntStream;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;
@@ -57,6 +58,7 @@ import org.springframework.cloud.stream.config.BindingServiceProperties;
import org.springframework.cloud.stream.messaging.DirectWithAttributesChannel;
import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.context.annotation.Bean;
import org.springframework.core.codec.CodecException;
import org.springframework.integration.channel.AbstractMessageChannel;
import org.springframework.integration.channel.DirectChannel;
import org.springframework.integration.config.GlobalChannelInterceptor;
@@ -203,6 +205,23 @@ class StreamBridgeTests {
}
}
// For more context on this test: https://github.com/spring-cloud/spring-cloud-stream/issues/3078
@Test
void functionInvocationWrapperNullError() {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(
EmptyConfiguration.class)).web(WebApplicationType.NONE).run(
"--spring.cloud.stream.source=outputA",
"--spring.jmx.enabled=false")) {
StreamBridge streamBridge = context.getBean(StreamBridge.class);
var exception = Assertions.assertThrows(RuntimeException.class, () -> streamBridge.send("outputA-out-0",
new CodecException("invalidException")
));
assertThat(exception.getMessage()).isEqualTo("org.springframework.cloud.function.context.catalog.SimpleFunctionRegistry$FunctionInvocationWrapper returned null");
}
}
// For more context on this test: https://github.com/spring-cloud/spring-cloud-stream/issues/2815
@Test
void ensurePartitioningWorksWhenNativeEncodingEnabled() {

View File

@@ -228,9 +228,13 @@ public final class StreamBridge implements StreamOperations, SmartInitializingSi
lock.unlock();
}
if (resultMessage == null
&& ((Message) messageToSend).getPayload().getClass().getName().equals("org.springframework.kafka.support.KafkaNull")) {
resultMessage = messageToSend;
if (resultMessage == null) {
if (((Message) messageToSend).getPayload().getClass().getName().equals("org.springframework.kafka.support.KafkaNull")) {
resultMessage = messageToSend;
}
else {
throw new RuntimeException(functionToInvoke.getClass().getName() + " returned null");
}
}
resultMessage = (Message<?>) this.functionInvocationHelper.postProcessResult(resultMessage, null);