diff --git a/spring-cloud-stream-test-support/src/main/java/org/springframework/cloud/stream/test/matcher/MessageQueueMatcher.java b/spring-cloud-stream-test-support/src/main/java/org/springframework/cloud/stream/test/matcher/MessageQueueMatcher.java index 28b406835..b8b627ed3 100644 --- a/spring-cloud-stream-test-support/src/main/java/org/springframework/cloud/stream/test/matcher/MessageQueueMatcher.java +++ b/spring-cloud-stream-test-support/src/main/java/org/springframework/cloud/stream/test/matcher/MessageQueueMatcher.java @@ -1,5 +1,6 @@ /* * Copyright 2015 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. * You may obtain a copy of the License at @@ -70,11 +71,10 @@ public class MessageQueueMatcher extends BaseMatcher public MessageQueueMatcher(Matcher delegate, long timeout, TimeUnit unit, Extractor, T> extractor) { this.delegate = delegate; this.timeout = timeout; - this.unit = unit; + this.unit = (unit != null ? unit : TimeUnit.SECONDS); this.extractor = extractor; } - @Override public boolean matches(Object item) { @SuppressWarnings("unchecked") @@ -83,9 +83,11 @@ public class MessageQueueMatcher extends BaseMatcher try { if (timeout > 0) { received = queue.poll(timeout, unit); - } else if (timeout == 0) { + } + else if (timeout == 0) { received = queue.poll(); - } else { + } + else { received = queue.take(); } } @@ -104,7 +106,8 @@ public class MessageQueueMatcher extends BaseMatcher T value = actuallyReceived.get(queue); if (value != null) { description.appendText("received: ").appendValue(value); - } else { + } + else { description.appendText("timed out after " + timeout + " " + unit.name().toLowerCase()); } } @@ -126,9 +129,9 @@ public class MessageQueueMatcher extends BaseMatcher description.appendText("Channel to receive ").appendDescriptionOf(extractor).appendDescriptionOf(delegate); } - @SuppressWarnings("unchecked") + @SuppressWarnings({"unchecked", "rawtypes"}) public static

MessageQueueMatcher

receivesMessageThat(Matcher> messageMatcher) { - return new MessageQueueMatcher(messageMatcher, 0, null, new Extractor, Message

>("a message that ") { + return new MessageQueueMatcher(messageMatcher, 5, TimeUnit.SECONDS, new Extractor, Message

>("a message that ") { @Override public Message

apply(Message

m) { return m; @@ -136,9 +139,9 @@ public class MessageQueueMatcher extends BaseMatcher }); } - @SuppressWarnings("unchecked") + @SuppressWarnings({"unchecked", "rawtypes"}) public static

MessageQueueMatcher

receivesPayloadThat(Matcher

payloadMatcher) { - return new MessageQueueMatcher(payloadMatcher, 0, null, new Extractor, P>("a message whose payload ") { + return new MessageQueueMatcher(payloadMatcher, 5, TimeUnit.SECONDS, new Extractor, P>("a message whose payload ") { @Override public P apply(Message

m) { return m.getPayload();