GH-295 Fixed tests to comply with double-conversion related changes in core
See https://github.com/spring-cloud/spring-cloud-stream/issues/1130 See https://github.com/spring-cloud/spring-cloud-stream/issues/1071 Resolves #295
This commit is contained in:
@@ -315,10 +315,10 @@ public class KafkaBinderTests extends
|
||||
|
||||
moduleOutputChannel.send(message);
|
||||
CountDownLatch latch = new CountDownLatch(1);
|
||||
AtomicReference<Message<String>> inboundMessageRef = new AtomicReference<>();
|
||||
AtomicReference<Message<byte[]>> inboundMessageRef = new AtomicReference<>();
|
||||
moduleInputChannel.subscribe(message1 -> {
|
||||
try {
|
||||
inboundMessageRef.set((Message<String>) message1);
|
||||
inboundMessageRef.set((Message<byte[]>) message1);
|
||||
}
|
||||
finally {
|
||||
latch.countDown();
|
||||
@@ -328,7 +328,7 @@ public class KafkaBinderTests extends
|
||||
|
||||
|
||||
Assertions.assertThat(inboundMessageRef.get()).isNotNull();
|
||||
Assertions.assertThat(inboundMessageRef.get().getPayload()).isEqualTo("foo");
|
||||
Assertions.assertThat(new String(inboundMessageRef.get().getPayload(), StandardCharsets.UTF_8)).isEqualTo("foo");
|
||||
Assertions.assertThat(inboundMessageRef.get().getHeaders().get(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE)).isNull();
|
||||
Assertions.assertThat(inboundMessageRef.get().getHeaders().get(MessageHeaders.CONTENT_TYPE))
|
||||
.isEqualTo(MimeTypeUtils.TEXT_PLAIN);
|
||||
@@ -366,10 +366,10 @@ public class KafkaBinderTests extends
|
||||
.setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN).build();
|
||||
moduleOutputChannel.send(message);
|
||||
CountDownLatch latch = new CountDownLatch(1);
|
||||
AtomicReference<Message<String>> inboundMessageRef = new AtomicReference<>();
|
||||
AtomicReference<Message<byte[]>> inboundMessageRef = new AtomicReference<>();
|
||||
moduleInputChannel.subscribe(message1 -> {
|
||||
try {
|
||||
inboundMessageRef.set((Message<String>) message1);
|
||||
inboundMessageRef.set((Message<byte[]>) message1);
|
||||
}
|
||||
finally {
|
||||
latch.countDown();
|
||||
@@ -378,7 +378,7 @@ public class KafkaBinderTests extends
|
||||
Assert.isTrue(latch.await(5, TimeUnit.SECONDS), "Failed to receive message");
|
||||
|
||||
assertThat(inboundMessageRef.get()).isNotNull();
|
||||
assertThat(inboundMessageRef.get().getPayload()).isEqualTo("foo");
|
||||
assertThat(new String(inboundMessageRef.get().getPayload(), StandardCharsets.UTF_8)).isEqualTo("foo");
|
||||
assertThat(inboundMessageRef.get().getHeaders().get(MessageHeaders.CONTENT_TYPE))
|
||||
.isEqualTo(MimeTypeUtils.TEXT_PLAIN);
|
||||
producerBinding.unbind();
|
||||
@@ -401,8 +401,7 @@ public class KafkaBinderTests extends
|
||||
moduleOutputChannel, outputBindingProperties.getProducer());
|
||||
Binding<MessageChannel> consumerBinding = binder.bindConsumer("foo.bar",
|
||||
"testSendAndReceive", moduleInputChannel, consumerProperties);
|
||||
// Bypass conversion we are only testing sendReceive
|
||||
Message<?> message = org.springframework.integration.support.MessageBuilder.withPayload("foo")
|
||||
Message<?> message = org.springframework.integration.support.MessageBuilder.withPayload("foo".getBytes(StandardCharsets.UTF_8))
|
||||
.setHeader(MessageHeaders.CONTENT_TYPE,
|
||||
MimeTypeUtils.APPLICATION_OCTET_STREAM)
|
||||
.build();
|
||||
@@ -410,10 +409,10 @@ public class KafkaBinderTests extends
|
||||
binderBindUnbindLatency();
|
||||
moduleOutputChannel.send(message);
|
||||
CountDownLatch latch = new CountDownLatch(1);
|
||||
AtomicReference<Message<String>> inboundMessageRef = new AtomicReference<>();
|
||||
AtomicReference<Message<byte[]>> inboundMessageRef = new AtomicReference<>();
|
||||
moduleInputChannel.subscribe(message1 -> {
|
||||
try {
|
||||
inboundMessageRef.set((Message<String>) message1);
|
||||
inboundMessageRef.set((Message<byte[]>) message1);
|
||||
}
|
||||
finally {
|
||||
latch.countDown();
|
||||
@@ -422,7 +421,7 @@ public class KafkaBinderTests extends
|
||||
Assert.isTrue(latch.await(5, TimeUnit.SECONDS), "Failed to receive message");
|
||||
|
||||
assertThat(inboundMessageRef.get()).isNotNull();
|
||||
assertThat(inboundMessageRef.get().getPayload()).isEqualTo("foo");
|
||||
assertThat(new String(inboundMessageRef.get().getPayload(), StandardCharsets.UTF_8)).isEqualTo("foo");
|
||||
assertThat(inboundMessageRef.get().getHeaders().get(MessageHeaders.CONTENT_TYPE))
|
||||
.isEqualTo(MimeTypeUtils.APPLICATION_OCTET_STREAM);
|
||||
producerBinding.unbind();
|
||||
@@ -696,7 +695,7 @@ public class KafkaBinderTests extends
|
||||
assertThat(handler.getReceivedMessages().entrySet()).hasSize(1);
|
||||
Message<?> receivedMessage = handler.getReceivedMessages().entrySet().iterator().next().getValue();
|
||||
assertThat(receivedMessage).isNotNull();
|
||||
assertThat(receivedMessage.getPayload()).isEqualTo(testMessagePayload);
|
||||
assertThat(new String((byte[])receivedMessage.getPayload(), StandardCharsets.UTF_8)).isEqualTo(testMessagePayload);
|
||||
assertThat(handler.getInvocationCount()).isEqualTo(consumerProperties.getMaxAttempts());
|
||||
consumerBinding.unbind();
|
||||
|
||||
@@ -762,7 +761,7 @@ public class KafkaBinderTests extends
|
||||
assertThat(handler.getReceivedMessages().entrySet()).hasSize(1);
|
||||
Message<?> handledMessage = handler.getReceivedMessages().entrySet().iterator().next().getValue();
|
||||
assertThat(handledMessage).isNotNull();
|
||||
assertThat(handledMessage.getPayload()).isEqualTo(testMessagePayload);
|
||||
assertThat(new String((byte[])handledMessage.getPayload(), StandardCharsets.UTF_8)).isEqualTo(testMessagePayload);
|
||||
assertThat(handler.getInvocationCount()).isEqualTo(consumerProperties.getMaxAttempts());
|
||||
binderBindUnbindLatency();
|
||||
dlqConsumerBinding.unbind();
|
||||
@@ -832,7 +831,7 @@ public class KafkaBinderTests extends
|
||||
assertThat(handler.getReceivedMessages().entrySet()).hasSize(1);
|
||||
Message<?> handledMessage = handler.getReceivedMessages().entrySet().iterator().next().getValue();
|
||||
assertThat(handledMessage).isNotNull();
|
||||
assertThat(handledMessage.getPayload()).isEqualTo(testMessagePayload);
|
||||
assertThat(new String((byte[])handledMessage.getPayload(), StandardCharsets.UTF_8)).isEqualTo(testMessagePayload);
|
||||
assertThat(handler.getInvocationCount()).isEqualTo(consumerProperties.getMaxAttempts());
|
||||
binderBindUnbindLatency();
|
||||
dlqConsumerBinding.unbind();
|
||||
@@ -892,10 +891,10 @@ public class KafkaBinderTests extends
|
||||
moduleOutputChannel.send(message);
|
||||
CountDownLatch latch = new CountDownLatch(1);
|
||||
|
||||
AtomicReference<Message<String>> inboundMessageRef = new AtomicReference<>();
|
||||
AtomicReference<Message<byte[]>> inboundMessageRef = new AtomicReference<>();
|
||||
moduleInputChannel.subscribe(message1 -> {
|
||||
try {
|
||||
inboundMessageRef.set((Message<String>) message1);
|
||||
inboundMessageRef.set((Message<byte[]>) message1);
|
||||
}
|
||||
finally {
|
||||
latch.countDown();
|
||||
@@ -904,7 +903,7 @@ public class KafkaBinderTests extends
|
||||
Assert.isTrue(latch.await(5, TimeUnit.SECONDS), "Failed to receive message");
|
||||
|
||||
assertThat(inboundMessageRef.get()).isNotNull();
|
||||
assertThat(inboundMessageRef.get().getPayload().getBytes()).containsExactly(testPayload);
|
||||
assertThat(inboundMessageRef.get().getPayload()).containsExactly(testPayload);
|
||||
producerBinding.unbind();
|
||||
consumerBinding.unbind();
|
||||
}
|
||||
@@ -952,10 +951,10 @@ public class KafkaBinderTests extends
|
||||
output.send(new GenericMessage<>(testPayload2.getBytes()));
|
||||
|
||||
CountDownLatch latch1 = new CountDownLatch(1);
|
||||
AtomicReference<Message<String>> inboundMessageRef2 = new AtomicReference<>();
|
||||
AtomicReference<Message<byte[]>> inboundMessageRef2 = new AtomicReference<>();
|
||||
input1.subscribe(message1 -> {
|
||||
try {
|
||||
inboundMessageRef2.set((Message<String>) message1);
|
||||
inboundMessageRef2.set((Message<byte[]>) message1);
|
||||
}
|
||||
finally {
|
||||
latch1.countDown();
|
||||
@@ -964,7 +963,7 @@ public class KafkaBinderTests extends
|
||||
Assert.isTrue(latch1.await(5, TimeUnit.SECONDS), "Failed to receive message");
|
||||
|
||||
assertThat(inboundMessageRef2.get()).isNotNull();
|
||||
assertThat(inboundMessageRef2.get().getPayload()).isEqualTo(testPayload2);
|
||||
assertThat(new String(inboundMessageRef2.get().getPayload(), StandardCharsets.UTF_8)).isEqualTo(testPayload2);
|
||||
Thread.sleep(2000);
|
||||
producerBinding.unbind();
|
||||
consumerBinding.unbind();
|
||||
@@ -1515,10 +1514,10 @@ public class KafkaBinderTests extends
|
||||
moduleOutputChannel.send(message);
|
||||
|
||||
CountDownLatch latch = new CountDownLatch(1);
|
||||
AtomicReference<Message<String>> inboundMessageRef = new AtomicReference<>();
|
||||
AtomicReference<Message<byte[]>> inboundMessageRef = new AtomicReference<>();
|
||||
moduleInputChannel.subscribe(message1 -> {
|
||||
try {
|
||||
inboundMessageRef.set((Message<String>) message1);
|
||||
inboundMessageRef.set((Message<byte[]>) message1);
|
||||
}
|
||||
finally {
|
||||
latch.countDown();
|
||||
@@ -1527,7 +1526,7 @@ public class KafkaBinderTests extends
|
||||
Assert.isTrue(latch.await(5, TimeUnit.SECONDS), "Failed to receive message");
|
||||
|
||||
assertThat(inboundMessageRef.get()).isNotNull();
|
||||
assertThat(inboundMessageRef.get().getPayload()).isEqualTo("testSendAndReceiveWithRawMode");
|
||||
assertThat(new String(inboundMessageRef.get().getPayload(), StandardCharsets.UTF_8)).isEqualTo("testSendAndReceiveWithRawMode");
|
||||
producerBinding.unbind();
|
||||
consumerBinding.unbind();
|
||||
}
|
||||
@@ -1763,10 +1762,10 @@ public class KafkaBinderTests extends
|
||||
consumerProperties);
|
||||
|
||||
CountDownLatch latch = new CountDownLatch(1);
|
||||
AtomicReference<Message<String>> inboundMessageRef1 = new AtomicReference<>();
|
||||
AtomicReference<Message<byte[]>> inboundMessageRef1 = new AtomicReference<>();
|
||||
MessageHandler messageHandler = message1 -> {
|
||||
try {
|
||||
inboundMessageRef1.set((Message<String>) message1);
|
||||
inboundMessageRef1.set((Message<byte[]>) message1);
|
||||
}
|
||||
finally {
|
||||
latch.countDown();
|
||||
@@ -1775,17 +1774,17 @@ public class KafkaBinderTests extends
|
||||
input1.subscribe(messageHandler);
|
||||
Assert.isTrue(latch.await(5, TimeUnit.SECONDS), "Failed to receive message");
|
||||
assertThat(inboundMessageRef1.get()).isNotNull();
|
||||
assertThat(inboundMessageRef1.get().getPayload()).isEqualTo(testPayload1);
|
||||
assertThat(new String(inboundMessageRef1.get().getPayload(), StandardCharsets.UTF_8)).isEqualTo(testPayload1);
|
||||
|
||||
String testPayload2 = "foo-" + UUID.randomUUID().toString();
|
||||
input1.unsubscribe(messageHandler);
|
||||
output.send(new GenericMessage<>(testPayload2.getBytes()));
|
||||
|
||||
CountDownLatch latch1 = new CountDownLatch(1);
|
||||
AtomicReference<Message<String>> inboundMessageRef2 = new AtomicReference<>();
|
||||
AtomicReference<Message<byte[]>> inboundMessageRef2 = new AtomicReference<>();
|
||||
input1.subscribe(message1 -> {
|
||||
try {
|
||||
inboundMessageRef2.set((Message<String>) message1);
|
||||
inboundMessageRef2.set((Message<byte[]>) message1);
|
||||
}
|
||||
finally {
|
||||
latch1.countDown();
|
||||
@@ -1794,7 +1793,7 @@ public class KafkaBinderTests extends
|
||||
Assert.isTrue(latch1.await(5, TimeUnit.SECONDS), "Failed to receive message");
|
||||
|
||||
assertThat(inboundMessageRef2.get()).isNotNull();
|
||||
assertThat(inboundMessageRef2.get().getPayload()).isEqualTo(testPayload2);
|
||||
assertThat(new String(inboundMessageRef2.get().getPayload(), StandardCharsets.UTF_8)).isEqualTo(testPayload2);
|
||||
|
||||
producerBinding.unbind();
|
||||
consumerBinding.unbind();
|
||||
@@ -1822,10 +1821,10 @@ public class KafkaBinderTests extends
|
||||
consumerBinding = binder.bindConsumer(testTopicName, "startOffsets", input1,
|
||||
firstConsumerProperties);
|
||||
CountDownLatch latch = new CountDownLatch(1);
|
||||
AtomicReference<Message<String>> inboundMessageRef1 = new AtomicReference<>();
|
||||
AtomicReference<Message<byte[]>> inboundMessageRef1 = new AtomicReference<>();
|
||||
MessageHandler messageHandler = message1 -> {
|
||||
try {
|
||||
inboundMessageRef1.set((Message<String>) message1);
|
||||
inboundMessageRef1.set((Message<byte[]>) message1);
|
||||
}
|
||||
finally {
|
||||
latch.countDown();
|
||||
@@ -1840,10 +1839,10 @@ public class KafkaBinderTests extends
|
||||
assertThat(inboundMessageRef1.get().getPayload()).isNotNull();
|
||||
input1.unsubscribe(messageHandler);
|
||||
CountDownLatch latch1 = new CountDownLatch(1);
|
||||
AtomicReference<Message<String>> inboundMessageRef2 = new AtomicReference<>();
|
||||
AtomicReference<Message<byte[]>> inboundMessageRef2 = new AtomicReference<>();
|
||||
MessageHandler messageHandler1 = message1 -> {
|
||||
try {
|
||||
inboundMessageRef2.set((Message<String>) message1);
|
||||
inboundMessageRef2.set((Message<byte[]>) message1);
|
||||
}
|
||||
finally {
|
||||
latch1.countDown();
|
||||
@@ -1862,10 +1861,10 @@ public class KafkaBinderTests extends
|
||||
consumerBinding = binder.bindConsumer(testTopicName, "startOffsets", input1, consumerProperties);
|
||||
input1.unsubscribe(messageHandler1);
|
||||
CountDownLatch latch2 = new CountDownLatch(1);
|
||||
AtomicReference<Message<String>> inboundMessageRef3 = new AtomicReference<>();
|
||||
AtomicReference<Message<byte[]>> inboundMessageRef3 = new AtomicReference<>();
|
||||
MessageHandler messageHandler2 = message1 -> {
|
||||
try {
|
||||
inboundMessageRef3.set((Message<String>) message1);
|
||||
inboundMessageRef3.set((Message< byte[]>) message1);
|
||||
}
|
||||
finally {
|
||||
latch2.countDown();
|
||||
@@ -1876,7 +1875,7 @@ public class KafkaBinderTests extends
|
||||
output.send(new GenericMessage<>(testPayload3.getBytes()));
|
||||
Assert.isTrue(latch2.await(15, TimeUnit.SECONDS), "Failed to receive message");
|
||||
assertThat(inboundMessageRef3.get()).isNotNull();
|
||||
assertThat(new String(inboundMessageRef3.get().getPayload())).isEqualTo(testPayload3);
|
||||
assertThat(new String(inboundMessageRef3.get().getPayload(), StandardCharsets.UTF_8)).isEqualTo(testPayload3);
|
||||
}
|
||||
finally {
|
||||
if (consumerBinding != null) {
|
||||
@@ -1927,10 +1926,10 @@ public class KafkaBinderTests extends
|
||||
|
||||
Binding<MessageChannel> consumerBinding = binder.bindConsumer(testTopicName, "test", input, consumerProperties);
|
||||
CountDownLatch latch = new CountDownLatch(1);
|
||||
AtomicReference<Message<String>> inboundMessageRef = new AtomicReference<>();
|
||||
AtomicReference<Message<byte[]>> inboundMessageRef = new AtomicReference<>();
|
||||
input.subscribe(message1 -> {
|
||||
try {
|
||||
inboundMessageRef.set((Message<String>) message1);
|
||||
inboundMessageRef.set((Message<byte[]>) message1);
|
||||
}
|
||||
finally {
|
||||
latch.countDown();
|
||||
@@ -1939,7 +1938,7 @@ public class KafkaBinderTests extends
|
||||
Assert.isTrue(latch.await(5, TimeUnit.SECONDS), "Failed to receive message");
|
||||
|
||||
assertThat(inboundMessageRef.get()).isNotNull();
|
||||
assertThat(inboundMessageRef.get().getPayload()).isEqualTo(testPayload);
|
||||
assertThat(new String(inboundMessageRef.get().getPayload(), StandardCharsets.UTF_8)).isEqualTo(testPayload);
|
||||
|
||||
producerBinding.unbind();
|
||||
consumerBinding.unbind();
|
||||
@@ -2293,10 +2292,10 @@ public class KafkaBinderTests extends
|
||||
binderBindUnbindLatency();
|
||||
moduleOutputChannel.send(message);
|
||||
CountDownLatch latch = new CountDownLatch(1);
|
||||
AtomicReference<Message<String>> inboundMessageRef = new AtomicReference<>();
|
||||
AtomicReference<Message<byte[]>> inboundMessageRef = new AtomicReference<>();
|
||||
moduleInputChannel.subscribe(message1 -> {
|
||||
try {
|
||||
inboundMessageRef.set((Message<String>) message1);
|
||||
inboundMessageRef.set((Message<byte[]>) message1);
|
||||
}
|
||||
finally {
|
||||
latch.countDown();
|
||||
@@ -2305,7 +2304,7 @@ public class KafkaBinderTests extends
|
||||
Assert.isTrue(latch.await(5, TimeUnit.SECONDS), "Failed to receive message");
|
||||
|
||||
assertThat(inboundMessageRef.get()).isNotNull();
|
||||
assertThat(inboundMessageRef.get().getPayload()).isEqualTo("test");
|
||||
assertThat(new String(inboundMessageRef.get().getPayload(), StandardCharsets.UTF_8)).isEqualTo("test");
|
||||
assertThat(inboundMessageRef.get().getHeaders()).containsEntry("contentType", MimeTypeUtils.TEXT_PLAIN);
|
||||
}
|
||||
finally {
|
||||
@@ -2367,17 +2366,17 @@ public class KafkaBinderTests extends
|
||||
moduleOutputChannel3.send(message);
|
||||
Message<?> inbound = receive(bridged, 10_000);
|
||||
assertThat(inbound).isNotNull();
|
||||
assertThat(inbound.getPayload()).isEqualTo("testSendAndReceiveWithMixedMode");
|
||||
assertThat(new String((byte[])inbound.getPayload(), StandardCharsets.UTF_8)).isEqualTo("testSendAndReceiveWithMixedMode");
|
||||
assertThat(inbound.getHeaders().get("foo")).isEqualTo("bar");
|
||||
assertThat(inbound.getHeaders().get(BinderHeaders.NATIVE_HEADERS_PRESENT)).isNull();
|
||||
inbound = receive(bridged);
|
||||
assertThat(inbound).isNotNull();
|
||||
assertThat(inbound.getPayload()).isEqualTo("testSendAndReceiveWithMixedMode");
|
||||
assertThat(new String((byte[])inbound.getPayload(), StandardCharsets.UTF_8)).isEqualTo("testSendAndReceiveWithMixedMode");
|
||||
assertThat(inbound.getHeaders().get("foo")).isEqualTo("bar");
|
||||
assertThat(inbound.getHeaders().get(BinderHeaders.NATIVE_HEADERS_PRESENT)).isEqualTo(Boolean.TRUE);
|
||||
inbound = receive(bridged);
|
||||
assertThat(inbound).isNotNull();
|
||||
assertThat(inbound.getPayload()).isEqualTo("testSendAndReceiveWithMixedMode");
|
||||
assertThat(new String((byte[])inbound.getPayload(), StandardCharsets.UTF_8)).isEqualTo("testSendAndReceiveWithMixedMode");
|
||||
assertThat(inbound.getHeaders().get("foo")).isNull();
|
||||
assertThat(inbound.getHeaders().get(BinderHeaders.NATIVE_HEADERS_PRESENT)).isNull();
|
||||
|
||||
|
||||
Reference in New Issue
Block a user