From 286ba44ff1dbf7469c699931742a29fea1c8dbce Mon Sep 17 00:00:00 2001 From: Eric Bottard Date: Fri, 18 Dec 2015 16:28:12 +0100 Subject: [PATCH] Remove TODOs --- .../binder/kafka/WindowingOffsetManager.java | 13 +++++---- .../stream/binder/kafka/KafkaBinderTests.java | 7 +++-- .../stream/binder/kafka/KafkaTestBinder.java | 7 +++-- .../rabbit/ConnectionFactorySettings.java | 1 - .../binder/rabbit/RabbitBindingCleaner.java | 1 - .../stream/binder/redis/RedisBinderTests.java | 6 ++--- .../MessageChannelBinderSupportTests.java | 6 ++--- .../module/launcher/ModuleLauncher.java | 1 - .../module/launcher/MultiArchiveLauncher.java | 2 -- .../tuple/DefaultTupleTestForBatch.java | 27 +++++++++---------- .../aggregation/ModuleAggregationTest.java | 3 +-- .../local/LocalMessageChannelBinder.java | 2 -- 12 files changed, 30 insertions(+), 46 deletions(-) diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/WindowingOffsetManager.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/WindowingOffsetManager.java index a82252787..e623b6f07 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/WindowingOffsetManager.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/WindowingOffsetManager.java @@ -21,12 +21,6 @@ import java.util.Collection; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; -import org.springframework.beans.factory.DisposableBean; -import org.springframework.beans.factory.InitializingBean; -import org.springframework.integration.kafka.core.Partition; -import org.springframework.integration.kafka.listener.OffsetManager; -import org.springframework.util.Assert; - import rx.Observable; import rx.Subscription; import rx.functions.Action0; @@ -39,6 +33,12 @@ import rx.subjects.PublishSubject; import rx.subjects.SerializedSubject; import rx.subjects.Subject; +import org.springframework.beans.factory.DisposableBean; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.integration.kafka.core.Partition; +import org.springframework.integration.kafka.listener.OffsetManager; +import org.springframework.util.Assert; + /** * An {@link OffsetManager} that aggregates writes over a time or count window, using an underlying delegate to * do the actual operations. Its purpose is to reduce the performance impact of writing operations @@ -48,7 +48,6 @@ import rx.subjects.Subject; * * @author Marius Bogoevici */ -//TODO: Move this class to spring-integration-kafka public class WindowingOffsetManager implements OffsetManager, InitializingBean, DisposableBean { private final CreatePartitionAndOffsetFunction createPartitionAndOffsetFunction = new CreatePartitionAndOffsetFunction(); diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java index a5d5ae290..249172a09 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java @@ -29,6 +29,7 @@ import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.BlockingQueue; import java.util.concurrent.TimeUnit; +import kafka.api.OffsetRequest; import org.junit.ClassRule; import org.junit.Ignore; import org.junit.Test; @@ -47,8 +48,6 @@ import org.springframework.integration.kafka.listener.MessageListener; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; -import kafka.api.OffsetRequest; - /** * Integration tests for the {@link KafkaMessageChannelBinder}. @@ -312,11 +311,11 @@ public class KafkaBinderTests extends PartitionCapableBinderTests { binder.unbindConsumers("foo" + uniqueBindingId + ".0"); } - @Override @Ignore // TODO + @Override @Ignore("https://github.com/spring-cloud/spring-cloud-stream/issues/243") public void testSendAndReceivePubSub() throws Exception { } - @Override @Ignore // TODO + @Override @Ignore("https://github.com/spring-cloud/spring-cloud-stream/issues/243") public void createInboundPubSubBeforeOutboundPubSub() throws Exception { } diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaTestBinder.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaTestBinder.java index c2d4a61c9..adf323726 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaTestBinder.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaTestBinder.java @@ -18,6 +18,9 @@ package org.springframework.cloud.stream.binder.kafka; import java.util.List; +import com.esotericsoftware.kryo.Kryo; +import com.esotericsoftware.kryo.Registration; + import org.springframework.cloud.stream.binder.AbstractTestBinder; import org.springframework.cloud.stream.test.junit.kafka.KafkaTestSupport; import org.springframework.cloud.stream.test.junit.kafka.TestKafkaCluster; @@ -28,9 +31,6 @@ import org.springframework.integration.codec.kryo.PojoCodec; import org.springframework.integration.kafka.support.ZookeeperConnect; import org.springframework.xd.tuple.serializer.kryo.TupleKryoRegistrar; -import com.esotericsoftware.kryo.Kryo; -import com.esotericsoftware.kryo.Registration; - /** * Test support class for {@link KafkaMessageChannelBinder}. @@ -79,7 +79,6 @@ public class KafkaTestBinder extends AbstractTestBinder 0); } @@ -332,8 +332,6 @@ public class DefaultTupleTestForBatch { } catch (ConversionFailedException e) { assertTrue(e.getMessage().indexOf("TestString") > 0); - // TODO - in batch this is part of the message, indicating what the name of the field is... - // assertTrue(e.getMessage().indexOf("name: [String]") > 0); } } @@ -382,7 +380,6 @@ public class DefaultTupleTestForBatch { } catch (ConversionFailedException e) { assertTrue(e.getMessage().indexOf("TestString") > 0); - // TODO - in batch this is part of the message, indicating what the name of the field is... // assertTrue(e.getMessage().indexOf("name: [String]") > 0); } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/aggregation/ModuleAggregationTest.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/aggregation/ModuleAggregationTest.java index ab221307a..ab3f25b9e 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/aggregation/ModuleAggregationTest.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/aggregation/ModuleAggregationTest.java @@ -35,8 +35,7 @@ import org.springframework.context.ConfigurableApplicationContext; /** * @author Marius Bogoevici */ -// TODO re-enable once we can test with a Mock binder -@Ignore +@Ignore("https://github.com/spring-cloud/spring-cloud-stream/issues/241") public class ModuleAggregationTest { @Test diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/local/LocalMessageChannelBinder.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/local/LocalMessageChannelBinder.java index 66503acbd..381353805 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/local/LocalMessageChannelBinder.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/local/LocalMessageChannelBinder.java @@ -272,7 +272,6 @@ public class LocalMessageChannelBinder extends MessageChannelBinderSupport { Properties properties) { validateConsumerProperties(name, properties, CONSUMER_REQUEST_REPLY_PROPERTIES); final MessageChannel requestChannel = this.findOrCreateRequestReplyChannel(name, "requestor.", properties); - // TODO: handle Pollable ? Assert.isInstanceOf(SubscribableChannel.class, requests); ((SubscribableChannel) requests).subscribe(new MessageHandler() { @@ -305,7 +304,6 @@ public class LocalMessageChannelBinder extends MessageChannelBinderSupport { } }); - // TODO: handle Pollable ? Assert.isInstanceOf(SubscribableChannel.class, replies); final SubscribableChannel replyChannel = this.findOrCreateRequestReplyChannel(name, "replier.", properties); ((SubscribableChannel) replies).subscribe(new MessageHandler() {