Fix Cyclomatic Complexity in Splitter

* Upgrade some dependencies
* Migrate `SplitterIntegrationTests` to JUnit 5
This commit is contained in:
Artem Bilan
2021-04-14 13:50:00 -04:00
parent 1a288cd527
commit e901c89fef
3 changed files with 143 additions and 157 deletions

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2019 the original author or authors.
* Copyright 2002-2021 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.
@@ -120,7 +120,6 @@ public abstract class AbstractMessageSplitter extends AbstractReplyProducingMess
}
@Override
@SuppressWarnings("unchecked")
protected final Object handleRequestMessage(Message<?> message) {
Object result = splitMessage(message);
// return null if 'null'
@@ -131,72 +130,90 @@ public abstract class AbstractMessageSplitter extends AbstractReplyProducingMess
boolean reactive = getOutputChannel() instanceof ReactiveStreamsSubscribableChannel;
setAsync(reactive);
Iterator<Object> iterator = null;
Flux<Object> flux = null;
final int sequenceSize;
if (reactive) {
return prepareFluxResult(message, result);
}
else {
return prepareIteratorResult(message, result);
}
}
@SuppressWarnings("unchecked")
private Flux<?> prepareFluxResult(Message<?> message, Object result) {
int sequenceSize = 1;
Flux<?> flux = Flux.just(result);
if (result instanceof Iterable<?>) {
Iterable<Object> iterable = (Iterable<Object>) result;
sequenceSize = obtainSizeIfPossible(iterable);
if (reactive) {
flux = Flux.fromIterable(iterable);
}
else {
iterator = iterable.iterator();
}
flux = Flux.fromIterable(iterable);
}
else if (result.getClass().isArray()) {
Object[] items = ObjectUtils.toObjectArray(result);
sequenceSize = items.length;
if (reactive) {
flux = Flux.fromArray(items);
}
else {
iterator = Arrays.asList(items).iterator();
}
flux = Flux.fromArray(items);
}
else if (result instanceof Iterator<?>) {
Iterator<Object> iter = (Iterator<Object>) result;
sequenceSize = obtainSizeIfPossible(iter);
if (reactive) {
flux = Flux.fromIterable(() -> iter);
}
else {
iterator = iter;
}
flux = Flux.fromIterable(() -> iter);
}
else if (result instanceof Stream<?>) {
Stream<Object> stream = ((Stream<Object>) result);
sequenceSize = 0;
if (reactive) {
flux = Flux.fromStream(stream);
}
else {
iterator = stream.iterator();
}
flux = Flux.fromStream(stream);
}
else if (result instanceof Publisher<?>) {
Publisher<Object> publisher = (Publisher<Object>) result;
sequenceSize = 0;
if (reactive) {
flux = Flux.from(publisher);
}
else {
iterator = Flux.from((Publisher<Object>) result).toIterable().iterator();
}
}
else {
sequenceSize = 1;
if (reactive) {
flux = Flux.just(result);
}
else {
iterator = Collections.singleton(result).iterator();
}
flux = Flux.from(publisher);
}
if (iterator != null && !iterator.hasNext()) {
Function<Object, ?> messageBuilderFunction = prepareMessageBuilderFunction(message, sequenceSize);
return flux
.map(messageBuilderFunction)
.switchIfEmpty(
Mono.defer(() -> {
MessageChannel discardingChannel = getDiscardChannel();
if (discardingChannel != null) {
this.messagingTemplate.send(discardingChannel, message);
}
return Mono.empty();
}));
}
@SuppressWarnings("unchecked")
private Iterator<?> prepareIteratorResult(Message<?> message, Object result) {
int sequenceSize = 1;
Iterator<?> iterator = Collections.singleton(result).iterator();
if (result instanceof Iterable<?>) {
Iterable<Object> iterable = (Iterable<Object>) result;
sequenceSize = obtainSizeIfPossible(iterable);
iterator = iterable.iterator();
}
else if (result.getClass().isArray()) {
Object[] items = ObjectUtils.toObjectArray(result);
sequenceSize = items.length;
iterator = Arrays.asList(items).iterator();
}
else if (result instanceof Iterator<?>) {
Iterator<Object> iter = (Iterator<Object>) result;
sequenceSize = obtainSizeIfPossible(iter);
iterator = iter;
}
else if (result instanceof Stream<?>) {
Stream<Object> stream = ((Stream<Object>) result);
sequenceSize = 0;
iterator = stream.iterator();
}
else if (result instanceof Publisher<?>) {
sequenceSize = 0;
iterator = Flux.from((Publisher<?>) result).toIterable().iterator();
}
if (!iterator.hasNext()) {
MessageChannel discardingChannel = getDiscardChannel();
if (discardingChannel != null) {
this.messagingTemplate.send(discardingChannel, message);
@@ -204,35 +221,25 @@ public abstract class AbstractMessageSplitter extends AbstractReplyProducingMess
return null;
}
Function<Object, ?> messageBuilderFunction = prepareMessageBuilderFunction(message, sequenceSize);
return new FunctionIterator<>(
result instanceof AutoCloseable && !result.equals(iterator) ? (AutoCloseable) result : null,
iterator, messageBuilderFunction);
}
private Function<Object, ?> prepareMessageBuilderFunction(Message<?> message, int sequenceSize) {
Map<String, Object> messageHeaders = message.getHeaders();
if (willAddHeaders(message)) {
messageHeaders = new HashMap<>(messageHeaders);
addHeaders(message, messageHeaders);
}
final Map<String, Object> headers = messageHeaders;
final Object correlationId = message.getHeaders().getId();
final AtomicInteger sequenceNumber = new AtomicInteger(1);
Map<String, Object> headers = messageHeaders;
Object correlationId = message.getHeaders().getId();
AtomicInteger sequenceNumber = new AtomicInteger(1);
Function<Object, AbstractIntegrationMessageBuilder<?>> messageBuilderFunction =
object -> createBuilder(object, headers, correlationId, sequenceNumber.getAndIncrement(), sequenceSize);
if (reactive) {
return flux
.map(messageBuilderFunction)
.switchIfEmpty(
Mono.defer(() -> {
MessageChannel discardingChannel = getDiscardChannel();
if (discardingChannel != null) {
this.messagingTemplate.send(discardingChannel, message);
}
return Mono.empty();
}));
}
else {
return new FunctionIterator<>(result instanceof AutoCloseable && !result.equals(iterator)
? (AutoCloseable) result : null, iterator, messageBuilderFunction);
}
return object -> createBuilder(object, headers, correlationId, sequenceNumber.getAndIncrement(), sequenceSize);
}
/**
@@ -349,7 +356,6 @@ public abstract class AbstractMessageSplitter extends AbstractReplyProducingMess
return object instanceof TreeNode;
}
@SuppressWarnings("unchecked")
private static int nodeSize(Object node) {
return ((TreeNode) node).size();
}

View File

@@ -17,15 +17,15 @@
package org.springframework.integration.splitter;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatExceptionOfType;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Iterator;
import java.util.List;
import org.junit.Before;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.BeanCreationException;
import org.springframework.beans.factory.annotation.Autowired;
@@ -39,19 +39,17 @@ import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.support.GenericMessage;
import org.springframework.test.annotation.DirtiesContext;
import org.springframework.test.annotation.DirtiesContext.ClassMode;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
import org.springframework.test.context.junit.jupiter.SpringJUnitConfig;
/**
* @author Iwein Fuld
* @author Alexander Peters
* @author Mark Fisher
* @author Gary Russell
* @author Artem Bilan
*/
@ContextConfiguration
@RunWith(SpringJUnit4ClassRunner.class)
@DirtiesContext(classMode = ClassMode.AFTER_EACH_TEST_METHOD)
@SpringJUnitConfig
@DirtiesContext
public class SplitterIntegrationTests {
@Autowired
@@ -77,19 +75,21 @@ public class SplitterIntegrationTests {
@Autowired
Receiver receiver;
@Before
@BeforeEach
public void clearWords() {
receiver.receivedWords.clear();
}
@MessageEndpoint
public static class Receiver {
private final List<String> receivedWords = new ArrayList<String>();
private final List<String> receivedWords = new ArrayList<>();
@ServiceActivator(inputChannel = "out")
public void deliveredWords(String string) {
this.receivedWords.add(string);
}
}
@MessageEndpoint
@@ -99,88 +99,68 @@ public class SplitterIntegrationTests {
public Iterator<String> split(String sentence) {
return Arrays.asList(sentence.split("\\s")).iterator();
}
}
@Test
public void configOk() throws Exception {
// just checking the parsing
}
@Test
public void annotated() throws Exception {
inAnnotated.send(new GenericMessage<String>(sentence));
inAnnotated.send(new GenericMessage<>(sentence));
assertThat(this.receiver.receivedWords.containsAll(words)).isTrue();
assertThat(words.containsAll(this.receiver.receivedWords)).isTrue();
}
@Test
public void methodInvoking() throws Exception {
inMethodInvoking.send(new GenericMessage<String>(sentence));
public void methodInvoking() {
inMethodInvoking.send(new GenericMessage<>(sentence));
assertThat(receiver.receivedWords.containsAll(words)).isTrue();
assertThat(words.containsAll(this.receiver.receivedWords)).isTrue();
}
@Test
public void defaultSplitter() throws Exception {
inDefault.send(new GenericMessage<List<String>>(words));
public void defaultSplitter() {
inDefault.send(new GenericMessage<>(words));
assertThat(receiver.receivedWords.containsAll(words)).isTrue();
assertThat(words.containsAll(receiver.receivedWords)).isTrue();
}
@Test
public void delimiterSplitter() throws Exception {
inDelimiters.send(new GenericMessage<String>("one,two, three; four/five"));
public void delimiterSplitter() {
inDelimiters.send(new GenericMessage<>("one,two, three; four/five"));
assertThat(receiver.receivedWords.containsAll(Arrays.asList("one", "two", "three", "four", "five"))).isTrue();
}
@Test(expected = IllegalArgumentException.class)
public void delimitersNotAllowedWithRef() throws Throwable {
try {
new ClassPathXmlApplicationContext("SplitterIntegrationTests-invalidRef.xml",
SplitterIntegrationTests.class).close();
}
catch (BeanCreationException e) {
Throwable cause = e.getMostSpecificCause();
assertThat(cause).isNotNull();
assertThat(cause instanceof IllegalArgumentException).isTrue();
assertThat(cause.getMessage().contains("'delimiters' property is only available")).isTrue();
throw cause;
}
}
@Test(expected = IllegalArgumentException.class)
public void delimitersNotAllowedWithInnerBean() throws Throwable {
try {
new ClassPathXmlApplicationContext("SplitterIntegrationTests-invalidInnerBean.xml",
SplitterIntegrationTests.class).close();
}
catch (BeanCreationException e) {
Throwable cause = e.getMostSpecificCause();
assertThat(cause).isNotNull();
assertThat(cause instanceof IllegalArgumentException).isTrue();
assertThat(cause.getMessage().contains("'delimiters' property is only available")).isTrue();
throw cause;
}
}
@Test(expected = IllegalArgumentException.class)
public void delimitersNotAllowedWithExpression() throws Throwable {
try {
new ClassPathXmlApplicationContext("SplitterIntegrationTests-invalidExpression.xml",
SplitterIntegrationTests.class).close();
}
catch (BeanCreationException e) {
Throwable cause = e.getMostSpecificCause();
assertThat(cause).isNotNull();
assertThat(cause instanceof IllegalArgumentException).isTrue();
assertThat(cause.getMessage().contains("'delimiters' property is only available")).isTrue();
throw cause;
}
@Test
public void delimitersNotAllowedWithRef() {
assertThatExceptionOfType(BeanCreationException.class)
.isThrownBy(() ->
new ClassPathXmlApplicationContext("SplitterIntegrationTests-invalidRef.xml",
SplitterIntegrationTests.class))
.withRootCauseExactlyInstanceOf(IllegalArgumentException.class)
.withMessageContaining("'delimiters' property is only available");
}
@Test
public void channelResolver_isNotNull() throws Exception {
public void delimitersNotAllowedWithInnerBean() {
assertThatExceptionOfType(BeanCreationException.class)
.isThrownBy(() ->
new ClassPathXmlApplicationContext("SplitterIntegrationTests-invalidInnerBean.xml",
SplitterIntegrationTests.class))
.withRootCauseExactlyInstanceOf(IllegalArgumentException.class)
.withMessageContaining("'delimiters' property is only available");
}
@Test
public void delimitersNotAllowedWithExpression() {
assertThatExceptionOfType(BeanCreationException.class)
.isThrownBy(() ->
new ClassPathXmlApplicationContext("SplitterIntegrationTests-invalidExpression.xml",
SplitterIntegrationTests.class))
.withRootCauseExactlyInstanceOf(IllegalArgumentException.class)
.withMessageContaining("'delimiters' property is only available");
}
@Test
public void channelResolver_isNotNull() {
splitter.setOutputChannel(null);
Message<String> message = MessageBuilder.withPayload("fooBar")
.setReplyChannelName("out").build();