diff --git a/org.springframework.integration.adapter/src/test/java/org/springframework/integration/adapter/stream/ByteStreamSourceTests.java b/org.springframework.integration.adapter/src/test/java/org/springframework/integration/adapter/stream/ByteStreamSourceTests.java index 69be88d7fb..e380bf921a 100644 --- a/org.springframework.integration.adapter/src/test/java/org/springframework/integration/adapter/stream/ByteStreamSourceTests.java +++ b/org.springframework.integration.adapter/src/test/java/org/springframework/integration/adapter/stream/ByteStreamSourceTests.java @@ -22,12 +22,10 @@ import static org.junit.Assert.assertNull; import java.io.ByteArrayInputStream; import org.junit.Test; - import org.springframework.integration.channel.MessageChannel; import org.springframework.integration.channel.QueueChannel; -import org.springframework.integration.endpoint.EndpointPoller; import org.springframework.integration.endpoint.SourceEndpoint; -import org.springframework.integration.message.GenericMessage; +import org.springframework.integration.endpoint.TriggerMessage; import org.springframework.integration.message.Message; import org.springframework.integration.scheduling.PollingSchedule; @@ -45,7 +43,7 @@ public class ByteStreamSourceTests { SourceEndpoint endpoint = new SourceEndpoint(source); endpoint.setTarget(channel); endpoint.afterPropertiesSet(); - endpoint.send(new GenericMessage(new EndpointPoller())); + endpoint.send(new TriggerMessage()); Message message1 = channel.receive(500); byte[] payload = (byte[]) message1.getPayload(); assertEquals(3, payload.length); @@ -54,7 +52,7 @@ public class ByteStreamSourceTests { assertEquals(3, payload[2]); Message message2 = channel.receive(0); assertNull(message2); - endpoint.send(new GenericMessage(new EndpointPoller())); + endpoint.send(new TriggerMessage()); Message message3 = channel.receive(0); assertNull(message3); } @@ -71,12 +69,12 @@ public class ByteStreamSourceTests { SourceEndpoint endpoint = new SourceEndpoint(source); endpoint.setTarget(channel); endpoint.afterPropertiesSet(); - endpoint.send(new GenericMessage(new EndpointPoller())); + endpoint.send(new TriggerMessage()); Message message1 = channel.receive(0); assertEquals(4, ((byte[]) message1.getPayload()).length); Message message2 = channel.receive(0); assertNull(message2); - endpoint.send(new GenericMessage(new EndpointPoller())); + endpoint.send(new TriggerMessage()); Message message3 = channel.receive(0); assertEquals(2, ((byte[]) message3.getPayload()).length); } @@ -94,12 +92,12 @@ public class ByteStreamSourceTests { SourceEndpoint endpoint = new SourceEndpoint(source); endpoint.setTarget(channel); endpoint.afterPropertiesSet(); - endpoint.send(new GenericMessage(new EndpointPoller())); + endpoint.send(new TriggerMessage()); Message message1 = channel.receive(0); assertEquals(4, ((byte[]) message1.getPayload()).length); Message message2 = channel.receive(0); assertNull(message2); - endpoint.send(new GenericMessage(new EndpointPoller())); + endpoint.send(new TriggerMessage()); Message message3 = channel.receive(0); assertEquals(4, ((byte[]) message3.getPayload()).length); assertEquals(0, ((byte[]) message3.getPayload())[3]); diff --git a/org.springframework.integration.adapter/src/test/java/org/springframework/integration/adapter/stream/CharacterStreamSourceTests.java b/org.springframework.integration.adapter/src/test/java/org/springframework/integration/adapter/stream/CharacterStreamSourceTests.java index 6a370adbed..fa387e914e 100644 --- a/org.springframework.integration.adapter/src/test/java/org/springframework/integration/adapter/stream/CharacterStreamSourceTests.java +++ b/org.springframework.integration.adapter/src/test/java/org/springframework/integration/adapter/stream/CharacterStreamSourceTests.java @@ -22,12 +22,10 @@ import static org.junit.Assert.assertNull; import java.io.StringReader; import org.junit.Test; - import org.springframework.integration.channel.MessageChannel; import org.springframework.integration.channel.QueueChannel; -import org.springframework.integration.endpoint.EndpointPoller; import org.springframework.integration.endpoint.SourceEndpoint; -import org.springframework.integration.message.GenericMessage; +import org.springframework.integration.endpoint.TriggerMessage; import org.springframework.integration.message.Message; import org.springframework.integration.scheduling.PollingSchedule; @@ -46,12 +44,12 @@ public class CharacterStreamSourceTests { SourceEndpoint endpoint = new SourceEndpoint(source); endpoint.setTarget(channel); endpoint.afterPropertiesSet(); - endpoint.send(new GenericMessage(new EndpointPoller())); + endpoint.send(new TriggerMessage()); Message message1 = channel.receive(0); assertEquals("test", message1.getPayload()); Message message2 = channel.receive(0); assertNull(message2); - endpoint.send(new GenericMessage(new EndpointPoller())); + endpoint.send(new TriggerMessage()); Message message3 = channel.receive(0); assertNull(message3); } diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/EndpointTrigger.java b/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/EndpointTrigger.java index a68d95d22e..287c603843 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/EndpointTrigger.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/EndpointTrigger.java @@ -18,7 +18,6 @@ package org.springframework.integration.endpoint; import org.springframework.integration.dispatcher.BroadcastingDispatcher; import org.springframework.integration.dispatcher.PollingDispatcher; -import org.springframework.integration.message.GenericMessage; import org.springframework.integration.message.Message; import org.springframework.integration.message.MessageSource; import org.springframework.integration.scheduling.PollingSchedule; @@ -36,7 +35,7 @@ public class EndpointTrigger extends PollingDispatcher { * Create an endpoint trigger with the specified {@link Schedule}. */ public EndpointTrigger(Schedule schedule) { - super(new EndpointPollerMessageSource(), new BroadcastingDispatcher(), schedule); + super(new TriggerSource(), new BroadcastingDispatcher(), schedule); } /** @@ -47,11 +46,19 @@ public class EndpointTrigger extends PollingDispatcher { this(new PollingSchedule(interval)); } + /** + * Create an endpoint trigger that will run one time only when submitted to + * a {@link org.springframework.integration.scheduling.TaskScheduler}. + */ + public EndpointTrigger() { + this(null); + } - private static class EndpointPollerMessageSource implements MessageSource { + + private static class TriggerSource implements MessageSource { public Message receive() { - return new GenericMessage(new EndpointPoller()); + return new TriggerMessage(); } } diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/TriggerMessage.java b/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/TriggerMessage.java new file mode 100644 index 0000000000..8bca61650d --- /dev/null +++ b/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/TriggerMessage.java @@ -0,0 +1,33 @@ +/* + * Copyright 2002-2008 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 + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.integration.endpoint; + +import org.springframework.integration.message.GenericMessage; + +/** + * A convenience Message implementation for sending a polling trigger + * to an endpoint. + * + * @author Mark Fisher + */ +public class TriggerMessage extends GenericMessage { + + public TriggerMessage() { + super(new EndpointPoller()); + } + +} diff --git a/org.springframework.integration/src/test/java/org/springframework/integration/config/EndpointInterceptorTests.java b/org.springframework.integration/src/test/java/org/springframework/integration/config/EndpointInterceptorTests.java index 33f57531ad..d940fbd409 100644 --- a/org.springframework.integration/src/test/java/org/springframework/integration/config/EndpointInterceptorTests.java +++ b/org.springframework.integration/src/test/java/org/springframework/integration/config/EndpointInterceptorTests.java @@ -21,15 +21,13 @@ import static org.junit.Assert.assertEquals; import java.util.List; import org.junit.Test; - import org.springframework.beans.DirectFieldAccessor; import org.springframework.context.support.ClassPathXmlApplicationContext; import org.springframework.integration.channel.MessageChannel; import org.springframework.integration.endpoint.EndpointInterceptor; -import org.springframework.integration.endpoint.EndpointPoller; import org.springframework.integration.endpoint.MessageEndpoint; import org.springframework.integration.endpoint.SourceEndpoint; -import org.springframework.integration.message.GenericMessage; +import org.springframework.integration.endpoint.TriggerMessage; import org.springframework.integration.message.StringMessage; /** @@ -105,7 +103,7 @@ public class EndpointInterceptorTests { if (endpoint instanceof SourceEndpoint) { MessageChannel channel = (MessageChannel) context.getBean("testChannel"); channel.send(new StringMessage("foo")); - endpoint.send(new GenericMessage(new EndpointPoller())); + endpoint.send(new TriggerMessage()); } else { endpoint.send(new StringMessage("test")); diff --git a/org.springframework.integration/src/test/java/org/springframework/integration/endpoint/SourceEndpointTests.java b/org.springframework.integration/src/test/java/org/springframework/integration/endpoint/SourceEndpointTests.java index c22bc4c150..bfba29075f 100644 --- a/org.springframework.integration/src/test/java/org/springframework/integration/endpoint/SourceEndpointTests.java +++ b/org.springframework.integration/src/test/java/org/springframework/integration/endpoint/SourceEndpointTests.java @@ -40,7 +40,7 @@ public class SourceEndpointTests { SourceEndpoint endpoint = new SourceEndpoint(source); endpoint.setTarget(channel); endpoint.afterPropertiesSet(); - endpoint.send(new GenericMessage(new EndpointPoller())); + endpoint.send(new TriggerMessage()); Message message = channel.receive(1000); assertNotNull("message should not be null", message); assertEquals("testing.1", message.getPayload());