From 8057b1f764d7419fdedfec95e0a46dcb6f5d4742 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Sat, 3 Feb 2018 07:48:44 -0500 Subject: [PATCH] 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 --- .../stream/binder/kafka/KafkaBinderTests.java | 89 +++++++++---------- 1 file changed, 44 insertions(+), 45 deletions(-) diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java index eedf59517..a5e2d45f7 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java @@ -315,10 +315,10 @@ public class KafkaBinderTests extends moduleOutputChannel.send(message); CountDownLatch latch = new CountDownLatch(1); - AtomicReference> inboundMessageRef = new AtomicReference<>(); + AtomicReference> inboundMessageRef = new AtomicReference<>(); moduleInputChannel.subscribe(message1 -> { try { - inboundMessageRef.set((Message) message1); + inboundMessageRef.set((Message) 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> inboundMessageRef = new AtomicReference<>(); + AtomicReference> inboundMessageRef = new AtomicReference<>(); moduleInputChannel.subscribe(message1 -> { try { - inboundMessageRef.set((Message) message1); + inboundMessageRef.set((Message) 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 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> inboundMessageRef = new AtomicReference<>(); + AtomicReference> inboundMessageRef = new AtomicReference<>(); moduleInputChannel.subscribe(message1 -> { try { - inboundMessageRef.set((Message) message1); + inboundMessageRef.set((Message) 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> inboundMessageRef = new AtomicReference<>(); + AtomicReference> inboundMessageRef = new AtomicReference<>(); moduleInputChannel.subscribe(message1 -> { try { - inboundMessageRef.set((Message) message1); + inboundMessageRef.set((Message) 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> inboundMessageRef2 = new AtomicReference<>(); + AtomicReference> inboundMessageRef2 = new AtomicReference<>(); input1.subscribe(message1 -> { try { - inboundMessageRef2.set((Message) message1); + inboundMessageRef2.set((Message) 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> inboundMessageRef = new AtomicReference<>(); + AtomicReference> inboundMessageRef = new AtomicReference<>(); moduleInputChannel.subscribe(message1 -> { try { - inboundMessageRef.set((Message) message1); + inboundMessageRef.set((Message) 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> inboundMessageRef1 = new AtomicReference<>(); + AtomicReference> inboundMessageRef1 = new AtomicReference<>(); MessageHandler messageHandler = message1 -> { try { - inboundMessageRef1.set((Message) message1); + inboundMessageRef1.set((Message) 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> inboundMessageRef2 = new AtomicReference<>(); + AtomicReference> inboundMessageRef2 = new AtomicReference<>(); input1.subscribe(message1 -> { try { - inboundMessageRef2.set((Message) message1); + inboundMessageRef2.set((Message) 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> inboundMessageRef1 = new AtomicReference<>(); + AtomicReference> inboundMessageRef1 = new AtomicReference<>(); MessageHandler messageHandler = message1 -> { try { - inboundMessageRef1.set((Message) message1); + inboundMessageRef1.set((Message) 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> inboundMessageRef2 = new AtomicReference<>(); + AtomicReference> inboundMessageRef2 = new AtomicReference<>(); MessageHandler messageHandler1 = message1 -> { try { - inboundMessageRef2.set((Message) message1); + inboundMessageRef2.set((Message) 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> inboundMessageRef3 = new AtomicReference<>(); + AtomicReference> inboundMessageRef3 = new AtomicReference<>(); MessageHandler messageHandler2 = message1 -> { try { - inboundMessageRef3.set((Message) 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 consumerBinding = binder.bindConsumer(testTopicName, "test", input, consumerProperties); CountDownLatch latch = new CountDownLatch(1); - AtomicReference> inboundMessageRef = new AtomicReference<>(); + AtomicReference> inboundMessageRef = new AtomicReference<>(); input.subscribe(message1 -> { try { - inboundMessageRef.set((Message) message1); + inboundMessageRef.set((Message) 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> inboundMessageRef = new AtomicReference<>(); + AtomicReference> inboundMessageRef = new AtomicReference<>(); moduleInputChannel.subscribe(message1 -> { try { - inboundMessageRef.set((Message) message1); + inboundMessageRef.set((Message) 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();