From 73b4b6fb3ebff3537bbf25691a9e11526423a9dc Mon Sep 17 00:00:00 2001 From: Marius Bogoevici Date: Fri, 19 Feb 2016 13:35:20 -0500 Subject: [PATCH] Fix race condition in multi-destination test --- .../binder/kafka/RawModeKafkaBinderTests.java | 35 ------------------ .../stream/binder/AbstractBinderTests.java | 37 +++++++++++-------- 2 files changed, 21 insertions(+), 51 deletions(-) diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/RawModeKafkaBinderTests.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/RawModeKafkaBinderTests.java index fa4e5d39a..962af27e9 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/RawModeKafkaBinderTests.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/RawModeKafkaBinderTests.java @@ -274,41 +274,6 @@ public class RawModeKafkaBinderTests extends KafkaBinderTests { assertTrue(getBindings(binder).isEmpty()); } - @Test - @Override - public void testSendAndReceiveMutipleTopics() throws Exception { - Binder binder = getBinder(); - - DirectChannel moduleOutputChannel1 = new DirectChannel(); - DirectChannel moduleOutputChannel2 = new DirectChannel(); - - QueueChannel moduleInputChannel = new QueueChannel(); - - Binding producerBinding1 = binder.bindProducer("foo.x", moduleOutputChannel1, null); - Binding producerBinding2 = binder.bindProducer("foo.y", moduleOutputChannel2, null); - - Binding consumerBinding1 = binder.bindConsumer("foo.x", "test", moduleInputChannel, null); - Binding consumerBinding2 = binder.bindConsumer("foo.y", "test", moduleInputChannel, null); - - Message message1 = MessageBuilder.withPayload("foo-x-payload".getBytes()).build(); - Message message2 = MessageBuilder.withPayload("foo-y-payload".getBytes()).build(); - - // Let the consumer actually bind to the producer before sending a msg - binderBindUnbindLatency(); - moduleOutputChannel1.send(message1); - Thread.sleep(50); - moduleOutputChannel2.send(message2); - - assertMessageReceive(moduleInputChannel, "foo-x-payload"); - assertMessageReceive(moduleInputChannel, "foo-y-payload"); - - binder.unbind(producerBinding1); - binder.unbind(consumerBinding1); - - binder.unbind(producerBinding2); - binder.unbind(consumerBinding2); - } - private void assertMessageReceive(QueueChannel moduleInputChannel, String payload) { Message inbound = receive(moduleInputChannel); assertNotNull(inbound); diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/AbstractBinderTests.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/AbstractBinderTests.java index b9c75f88b..499a67f52 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/AbstractBinderTests.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/AbstractBinderTests.java @@ -16,14 +16,19 @@ package org.springframework.cloud.stream.binder; +import static org.hamcrest.Matchers.equalTo; +import static org.hamcrest.Matchers.hasProperty; +import static org.hamcrest.collection.IsArrayContainingInAnyOrder.arrayContainingInAnyOrder; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertThat; import static org.junit.Assert.assertTrue; import java.util.Collection; import java.util.Collections; import java.util.List; +import java.util.UUID; import org.junit.After; import org.junit.Test; @@ -107,7 +112,7 @@ public abstract class AbstractBinderTests { } @Test - public void testSendAndReceiveMutipleTopics() throws Exception { + public void testSendAndReceiveMultipleTopics() throws Exception { Binder binder = getBinder(); DirectChannel moduleOutputChannel1 = new DirectChannel(); @@ -121,19 +126,27 @@ public abstract class AbstractBinderTests { Binding consumerBinding1 = binder.bindConsumer("foo.x", "test", moduleInputChannel, null); Binding consumerBinding2 = binder.bindConsumer("foo.y", "test", moduleInputChannel, null); - Message message1 = MessageBuilder.withPayload("foo-x-payload").setHeader(MessageHeaders.CONTENT_TYPE, - "foo/bar").build(); - Message message2 = MessageBuilder.withPayload("foo-y-payload").setHeader(MessageHeaders.CONTENT_TYPE, - "foo/bar").build(); + String testPayload1 = "foo" + UUID.randomUUID().toString(); + Message message1 = MessageBuilder.withPayload(testPayload1.getBytes()).build(); + String testPayload2 = "foo" + UUID.randomUUID().toString(); + Message message2 = MessageBuilder.withPayload(testPayload2.getBytes()).build(); // Let the consumer actually bind to the producer before sending a msg binderBindUnbindLatency(); moduleOutputChannel1.send(message1); - Thread.sleep(50); moduleOutputChannel2.send(message2); - assertMessageReceive(moduleInputChannel, "foo-x-payload"); - assertMessageReceive(moduleInputChannel, "foo-y-payload"); + + Message messages[] = new Message[2]; + messages[0] = receive(moduleInputChannel); + messages[1] = receive(moduleInputChannel); + + assertNotNull(messages[0]); + assertNotNull(messages[1]); + assertThat(messages, arrayContainingInAnyOrder( + hasProperty("payload", equalTo(testPayload1.getBytes())), + hasProperty("payload", equalTo(testPayload2.getBytes())))); + binder.unbind(producerBinding1); binder.unbind(consumerBinding1); @@ -142,14 +155,6 @@ public abstract class AbstractBinderTests { binder.unbind(consumerBinding2); } - private void assertMessageReceive(QueueChannel moduleInputChannel, String payload) { - Message inbound = receive(moduleInputChannel); - assertNotNull(inbound); - assertEquals(payload, inbound.getPayload()); - assertNull(inbound.getHeaders().get(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE)); - assertEquals("foo/bar", inbound.getHeaders().get(MessageHeaders.CONTENT_TYPE)); - } - @Test public void testSendAndReceiveNoOriginalContentType() throws Exception { Binder binder = getBinder();