Renamed PollingConsumerEndpoint to PollingConsumer.
This commit is contained in:
@@ -27,7 +27,7 @@ import org.springframework.integration.channel.PollableChannel;
|
||||
import org.springframework.integration.channel.SubscribableChannel;
|
||||
import org.springframework.integration.core.MessageChannel;
|
||||
import org.springframework.integration.endpoint.MessageEndpoint;
|
||||
import org.springframework.integration.endpoint.PollingConsumerEndpoint;
|
||||
import org.springframework.integration.endpoint.PollingConsumer;
|
||||
import org.springframework.integration.endpoint.SubscribingConsumerEndpoint;
|
||||
import org.springframework.integration.message.MessageHandler;
|
||||
import org.springframework.integration.scheduling.IntervalTrigger;
|
||||
@@ -152,7 +152,7 @@ public class ConsumerEndpointFactoryBean implements FactoryBean, BeanFactoryAwar
|
||||
if (this.trigger == null) {
|
||||
this.trigger = new IntervalTrigger(0);
|
||||
}
|
||||
PollingConsumerEndpoint pollingEndpoint = new PollingConsumerEndpoint(
|
||||
PollingConsumer pollingEndpoint = new PollingConsumer(
|
||||
(PollableChannel) channel, this.handler);
|
||||
pollingEndpoint.setTrigger(this.trigger);
|
||||
pollingEndpoint.setMaxMessagesPerPoll(this.maxMessagesPerPoll);
|
||||
|
||||
@@ -31,7 +31,7 @@ import org.springframework.integration.channel.PollableChannel;
|
||||
import org.springframework.integration.channel.SubscribableChannel;
|
||||
import org.springframework.integration.core.MessageChannel;
|
||||
import org.springframework.integration.endpoint.MessageEndpoint;
|
||||
import org.springframework.integration.endpoint.PollingConsumerEndpoint;
|
||||
import org.springframework.integration.endpoint.PollingConsumer;
|
||||
import org.springframework.integration.endpoint.SubscribingConsumerEndpoint;
|
||||
import org.springframework.integration.message.MessageHandler;
|
||||
import org.springframework.util.Assert;
|
||||
@@ -89,7 +89,7 @@ public abstract class AbstractMethodAnnotationPostProcessor<T extends Annotation
|
||||
MessageChannel inputChannel = this.channelResolver.resolveChannelName(inputChannelName);
|
||||
Assert.notNull(inputChannel, "failed to resolve inputChannel '" + inputChannelName + "'");
|
||||
if (inputChannel instanceof PollableChannel) {
|
||||
PollingConsumerEndpoint pollingEndpoint = new PollingConsumerEndpoint(
|
||||
PollingConsumer pollingEndpoint = new PollingConsumer(
|
||||
(PollableChannel) inputChannel, handler);
|
||||
if (pollerAnnotation != null) {
|
||||
AnnotationConfigUtils.configurePollingEndpointWithPollerAnnotation(
|
||||
|
||||
@@ -31,7 +31,7 @@ import org.springframework.integration.channel.SubscribableChannel;
|
||||
import org.springframework.integration.consumer.MethodInvokingMessageHandler;
|
||||
import org.springframework.integration.core.MessageChannel;
|
||||
import org.springframework.integration.endpoint.MessageEndpoint;
|
||||
import org.springframework.integration.endpoint.PollingConsumerEndpoint;
|
||||
import org.springframework.integration.endpoint.PollingConsumer;
|
||||
import org.springframework.integration.endpoint.SourcePollingChannelAdapter;
|
||||
import org.springframework.integration.endpoint.SubscribingConsumerEndpoint;
|
||||
import org.springframework.integration.message.MethodInvokingSource;
|
||||
@@ -115,7 +115,7 @@ public class ChannelAdapterAnnotationPostProcessor implements MethodAnnotationPo
|
||||
Trigger trigger = (pollerAnnotation != null)
|
||||
? AnnotationConfigUtils.parseTriggerFromPollerAnnotation(pollerAnnotation)
|
||||
: new IntervalTrigger(0);
|
||||
PollingConsumerEndpoint endpoint = new PollingConsumerEndpoint((PollableChannel) channel, handler);
|
||||
PollingConsumer endpoint = new PollingConsumer((PollableChannel) channel, handler);
|
||||
endpoint.setTrigger(trigger);
|
||||
return endpoint;
|
||||
}
|
||||
|
||||
@@ -22,9 +22,12 @@ import org.springframework.integration.message.MessageHandler;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
/**
|
||||
* Message Endpoint that connects any {@link MessageHandler} implementation
|
||||
* to a {@link PollableChannel}.
|
||||
*
|
||||
* @author Mark Fisher
|
||||
*/
|
||||
public class PollingConsumerEndpoint extends AbstractPollingEndpoint {
|
||||
public class PollingConsumer extends AbstractPollingEndpoint {
|
||||
|
||||
private final PollableChannel inputChannel;
|
||||
|
||||
@@ -33,7 +36,7 @@ public class PollingConsumerEndpoint extends AbstractPollingEndpoint {
|
||||
private volatile long receiveTimeout = 1000;
|
||||
|
||||
|
||||
public PollingConsumerEndpoint(PollableChannel inputChannel, MessageHandler handler) {
|
||||
public PollingConsumer(PollableChannel inputChannel, MessageHandler handler) {
|
||||
Assert.notNull(inputChannel, "inputChannel must not be null");
|
||||
Assert.notNull(handler, "handler must not be null");
|
||||
this.inputChannel = inputChannel;
|
||||
@@ -27,7 +27,7 @@ import org.springframework.integration.core.MessageChannel;
|
||||
import org.springframework.integration.core.MessagingException;
|
||||
import org.springframework.integration.endpoint.MessageEndpoint;
|
||||
import org.springframework.integration.endpoint.MessagingGateway;
|
||||
import org.springframework.integration.endpoint.PollingConsumerEndpoint;
|
||||
import org.springframework.integration.endpoint.PollingConsumer;
|
||||
import org.springframework.integration.endpoint.SubscribingConsumerEndpoint;
|
||||
import org.springframework.integration.message.ErrorMessage;
|
||||
import org.springframework.integration.message.MessageHandler;
|
||||
@@ -189,7 +189,7 @@ public abstract class AbstractMessagingGateway implements MessagingGateway, Mess
|
||||
(SubscribableChannel) this.replyChannel, handler);
|
||||
}
|
||||
else if (this.replyChannel instanceof PollableChannel) {
|
||||
PollingConsumerEndpoint endpoint = new PollingConsumerEndpoint(
|
||||
PollingConsumer endpoint = new PollingConsumer(
|
||||
(PollableChannel) this.replyChannel, handler);
|
||||
endpoint.setTaskScheduler(this.taskScheduler);
|
||||
endpoint.afterPropertiesSet();
|
||||
|
||||
@@ -40,7 +40,7 @@ import org.springframework.integration.config.xml.MessageBusParser;
|
||||
import org.springframework.integration.consumer.AbstractReplyProducingMessageHandler;
|
||||
import org.springframework.integration.consumer.ReplyMessageHolder;
|
||||
import org.springframework.integration.core.Message;
|
||||
import org.springframework.integration.endpoint.PollingConsumerEndpoint;
|
||||
import org.springframework.integration.endpoint.PollingConsumer;
|
||||
import org.springframework.integration.endpoint.SourcePollingChannelAdapter;
|
||||
import org.springframework.integration.endpoint.SubscribingConsumerEndpoint;
|
||||
import org.springframework.integration.message.ErrorMessage;
|
||||
@@ -75,7 +75,7 @@ public class ApplicationContextMessageBusTests {
|
||||
}
|
||||
};
|
||||
handler.setBeanFactory(context);
|
||||
PollingConsumerEndpoint endpoint = new PollingConsumerEndpoint(sourceChannel, handler);
|
||||
PollingConsumer endpoint = new PollingConsumer(sourceChannel, handler);
|
||||
endpoint.afterPropertiesSet();
|
||||
context.getBeanFactory().registerSingleton("testEndpoint", endpoint);
|
||||
context.refresh();
|
||||
@@ -148,9 +148,9 @@ public class ApplicationContextMessageBusTests {
|
||||
context.getBeanFactory().registerSingleton("output2", outputChannel2);
|
||||
handler1.setOutputChannel(outputChannel1);
|
||||
handler2.setOutputChannel(outputChannel2);
|
||||
PollingConsumerEndpoint endpoint1 = new PollingConsumerEndpoint(inputChannel, handler1);
|
||||
PollingConsumer endpoint1 = new PollingConsumer(inputChannel, handler1);
|
||||
endpoint1.afterPropertiesSet();
|
||||
PollingConsumerEndpoint endpoint2 = new PollingConsumerEndpoint(inputChannel, handler2);
|
||||
PollingConsumer endpoint2 = new PollingConsumer(inputChannel, handler2);
|
||||
endpoint2.afterPropertiesSet();
|
||||
context.getBeanFactory().registerSingleton("testEndpoint1", endpoint1);
|
||||
context.getBeanFactory().registerSingleton("testEndpoint2", endpoint2);
|
||||
@@ -266,7 +266,7 @@ public class ApplicationContextMessageBusTests {
|
||||
latch.countDown();
|
||||
}
|
||||
};
|
||||
PollingConsumerEndpoint endpoint = new PollingConsumerEndpoint(errorChannel, handler);
|
||||
PollingConsumer endpoint = new PollingConsumer(errorChannel, handler);
|
||||
endpoint.afterPropertiesSet();
|
||||
context.getBeanFactory().registerSingleton("testEndpoint", endpoint);
|
||||
ApplicationContextMessageBus bus = new ApplicationContextMessageBus();
|
||||
|
||||
@@ -17,9 +17,9 @@
|
||||
|
||||
<bean id="targetChannel" class="org.springframework.integration.channel.QueueChannel"/>
|
||||
|
||||
<bean id="endpoint" class="org.springframework.integration.endpoint.PollingConsumerEndpoint">
|
||||
<constructor-arg ref="serviceActivator"/>
|
||||
<bean id="endpoint" class="org.springframework.integration.endpoint.PollingConsumer">
|
||||
<constructor-arg ref="sourceChannel"/>
|
||||
<constructor-arg ref="serviceActivator"/>
|
||||
</bean>
|
||||
|
||||
<bean id="serviceActivator" class="org.springframework.integration.consumer.ServiceActivatingHandler">
|
||||
|
||||
@@ -35,7 +35,7 @@ import org.springframework.integration.consumer.AbstractReplyProducingMessageHan
|
||||
import org.springframework.integration.consumer.ReplyMessageHolder;
|
||||
import org.springframework.integration.core.Message;
|
||||
import org.springframework.integration.core.MessageChannel;
|
||||
import org.springframework.integration.endpoint.PollingConsumerEndpoint;
|
||||
import org.springframework.integration.endpoint.PollingConsumer;
|
||||
import org.springframework.integration.message.MessageBuilder;
|
||||
import org.springframework.integration.message.StringMessage;
|
||||
import org.springframework.integration.util.TestUtils;
|
||||
@@ -58,7 +58,7 @@ public class MessageChannelTemplateTests {
|
||||
replyHolder.set(message.getPayload().toString().toUpperCase());
|
||||
}
|
||||
};
|
||||
PollingConsumerEndpoint endpoint = new PollingConsumerEndpoint(requestChannel, handler);
|
||||
PollingConsumer endpoint = new PollingConsumer(requestChannel, handler);
|
||||
endpoint.afterPropertiesSet();
|
||||
GenericApplicationContext context = new GenericApplicationContext();
|
||||
context.getBeanFactory().registerSingleton("requestChannel", requestChannel);
|
||||
|
||||
@@ -46,7 +46,7 @@ import org.springframework.integration.channel.QueueChannel;
|
||||
import org.springframework.integration.config.xml.MessageBusParser;
|
||||
import org.springframework.integration.core.Message;
|
||||
import org.springframework.integration.core.MessageChannel;
|
||||
import org.springframework.integration.endpoint.PollingConsumerEndpoint;
|
||||
import org.springframework.integration.endpoint.PollingConsumer;
|
||||
import org.springframework.integration.message.MessageBuilder;
|
||||
import org.springframework.integration.message.MessageHandler;
|
||||
import org.springframework.integration.message.StringMessage;
|
||||
@@ -374,8 +374,7 @@ public class MessagingAnnotationPostProcessorTests {
|
||||
postProcessor.afterPropertiesSet();
|
||||
AnnotatedEndpointWithPolledAnnotation bean = new AnnotatedEndpointWithPolledAnnotation();
|
||||
postProcessor.postProcessAfterInitialization(bean, "testBean");
|
||||
PollingConsumerEndpoint endpoint =
|
||||
(PollingConsumerEndpoint) context.getBean("testBean.prependFoo.serviceActivator");
|
||||
PollingConsumer endpoint = (PollingConsumer) context.getBean("testBean.prependFoo.serviceActivator");
|
||||
Trigger trigger = (Trigger) new DirectFieldAccessor(endpoint).getPropertyValue("trigger");
|
||||
assertEquals(IntervalTrigger.class, trigger.getClass());
|
||||
DirectFieldAccessor triggerAccessor = new DirectFieldAccessor(trigger);
|
||||
|
||||
@@ -51,7 +51,7 @@ import org.springframework.integration.util.ErrorHandler;
|
||||
@SuppressWarnings("unchecked")
|
||||
public class PollingConsumerEndpointTests {
|
||||
|
||||
private PollingConsumerEndpoint endpoint;
|
||||
private PollingConsumer endpoint;
|
||||
|
||||
private TestTrigger trigger = new TestTrigger();
|
||||
|
||||
@@ -72,7 +72,7 @@ public class PollingConsumerEndpointTests {
|
||||
public void init() throws InterruptedException {
|
||||
consumer.counter.set(0);
|
||||
trigger.reset();
|
||||
endpoint = new PollingConsumerEndpoint(channelMock, consumer);
|
||||
endpoint = new PollingConsumer(channelMock, consumer);
|
||||
endpoint.setTaskScheduler(taskScheduler);
|
||||
taskScheduler.setErrorHandler(errorHandler);
|
||||
taskScheduler.start();
|
||||
|
||||
@@ -32,7 +32,7 @@ import org.springframework.integration.channel.QueueChannel;
|
||||
import org.springframework.integration.consumer.MethodInvokingMessageHandler;
|
||||
import org.springframework.integration.core.Message;
|
||||
import org.springframework.integration.core.MessagingException;
|
||||
import org.springframework.integration.endpoint.PollingConsumerEndpoint;
|
||||
import org.springframework.integration.endpoint.PollingConsumer;
|
||||
import org.springframework.integration.util.TestUtils;
|
||||
|
||||
/**
|
||||
@@ -85,7 +85,7 @@ public class MethodInvokingMessageHandlerTests {
|
||||
channel.send(message);
|
||||
assertNull(queue.poll());
|
||||
MethodInvokingMessageHandler handler = new MethodInvokingMessageHandler(testBean, "foo");
|
||||
PollingConsumerEndpoint endpoint = new PollingConsumerEndpoint(channel, handler);
|
||||
PollingConsumer endpoint = new PollingConsumer(channel, handler);
|
||||
context.getBeanFactory().registerSingleton("testEndpoint", endpoint);
|
||||
ApplicationContextMessageBus bus = new ApplicationContextMessageBus();
|
||||
bus.setTaskScheduler(TestUtils.createTaskScheduler(10));
|
||||
|
||||
Reference in New Issue
Block a user