diff --git a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/AbstractInboundTwitterEndpointSupport.java b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/AbstractInboundTwitterEndpointSupport.java index 92462ed722..bd08087349 100644 --- a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/AbstractInboundTwitterEndpointSupport.java +++ b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/AbstractInboundTwitterEndpointSupport.java @@ -20,9 +20,7 @@ import org.apache.commons.lang.exception.ExceptionUtils; import org.springframework.context.Lifecycle; import org.springframework.integration.Message; -import org.springframework.integration.MessageChannel; -import org.springframework.integration.core.MessagingTemplate; -import org.springframework.integration.endpoint.AbstractEndpoint; +import org.springframework.integration.endpoint.MessageProducerSupport; import org.springframework.integration.history.HistoryWritingMessagePostProcessor; import org.springframework.integration.history.TrackableComponent; import org.springframework.integration.support.MessageBuilder; @@ -36,6 +34,9 @@ import twitter4j.Twitter; import java.util.ArrayList; import java.util.List; +import java.util.concurrent.Executor; +import java.util.concurrent.Executors; +import java.util.concurrent.ThreadFactory; /** @@ -53,16 +54,35 @@ import java.util.List; * @author Josh Long * @since 2.0 */ -public abstract class AbstractInboundTwitterEndpointSupport extends AbstractEndpoint implements Lifecycle, TrackableComponent { +public abstract class AbstractInboundTwitterEndpointSupport extends MessageProducerSupport implements Lifecycle, TrackableComponent, Runnable { + protected volatile OAuthConfiguration configuration; - protected final MessagingTemplate messagingTemplate = new MessagingTemplate(); - private volatile MessageChannel requestChannel; protected volatile long markerId = -1; protected Twitter twitter; private final Object markerGuard = new Object(); private final Object apiPermitGuard = new Object(); - private final HistoryWritingMessagePostProcessor historyWritingPostProcessor = - new HistoryWritingMessagePostProcessor(); + private final HistoryWritingMessagePostProcessor historyWritingPostProcessor = new HistoryWritingMessagePostProcessor(); + protected Executor taskExecutor; + protected int poolSize = 1; + + public void setPoolSize(int poolSize) { + this.poolSize = poolSize; + } + + protected void checkTaskExecutor(final String threadName) { + if (this.taskExecutor == null) { + this.taskExecutor = Executors.newFixedThreadPool(this.poolSize, + new ThreadFactory() { + public Thread newThread(Runnable runner) { + Thread thread = new Thread(runner); + thread.setName(threadName); + thread.setDaemon(true); + + return thread; + } + }); + } + } public void setConfiguration(OAuthConfiguration configuration) { this.configuration = configuration; @@ -72,6 +92,14 @@ public abstract class AbstractInboundTwitterEndpointSupport extends AbstractE abstract protected List sort(List rl); + @Override + protected void onInit() { + super.onInit(); + Assert.notNull(this.configuration, "'configuration' can't be null"); + this.twitter = this.configuration.getTwitter(); + Assert.notNull(this.twitter, "'twitter' instance can't be null"); + } + protected void forwardAll(List tResponses) { List stats = new ArrayList(); @@ -86,35 +114,39 @@ public abstract class AbstractInboundTwitterEndpointSupport extends AbstractE return markerId; } - public String getComponentType() { - return "twitter:inbound-dm-channel-adapter"; - } - - public void setRequestChannel(MessageChannel requestChannel) { - this.messagingTemplate.setDefaultChannel(requestChannel); - this.requestChannel = requestChannel; - } + abstract public String getComponentType(); @Override protected void doStart() { - try { - this.historyWritingPostProcessor.setTrackableComponent(this); - refresh(); - } catch (Exception e) { - throw new RuntimeException(e); - } + historyWritingPostProcessor.setTrackableComponent(this); + + checkTaskExecutor(getClass().getName() + "-taskExecutor"); + + taskExecutor.execute(this); } protected void forward(T status) { synchronized (this.markerGuard) { Message twtMsg = MessageBuilder.withPayload(status).build(); - messagingTemplate.convertAndSend(requestChannel, twtMsg, - this.historyWritingPostProcessor); + + sendMessage(twtMsg); + markLastStatusId(status); } } - abstract protected void refresh() throws Exception; + /** + * this is execu + */ + public void run() { + try { + beginPolling(); + } catch (Exception e) { + throw new RuntimeException(e); + } + } + + abstract protected void beginPolling() throws Exception; protected void forwardAll(ResponseList tResponses) { List stats = new ArrayList(); @@ -127,7 +159,7 @@ public abstract class AbstractInboundTwitterEndpointSupport extends AbstractE } @SuppressWarnings("unchecked") - protected void runAsAPIRateLimitsPermit(ApiCallback cb) + protected void runAsAPIRateLimitsPermit(ApiCallback apiCallback) throws Exception { synchronized (this.apiPermitGuard) { while (waitUntilPullAvailable()) { @@ -135,7 +167,7 @@ public abstract class AbstractInboundTwitterEndpointSupport extends AbstractE logger.debug("have room to make an API request now"); } - cb.run(this, twitter); + apiCallback.run(this, twitter); } } } @@ -171,8 +203,8 @@ public abstract class AbstractInboundTwitterEndpointSupport extends AbstractE Thread.sleep(msUntilWeCanPullAgain); } catch (Throwable throwable) { - logger.debug("encountered an error when" + - " trying to refresh the timeline: " + + logger.debug( + "encountered an error when trying to refresh the timeline: " + ExceptionUtils.getFullStackTrace(throwable)); } @@ -187,14 +219,6 @@ public abstract class AbstractInboundTwitterEndpointSupport extends AbstractE return markerId > -1; } - @Override - protected void onInit() throws Exception { - messagingTemplate.afterPropertiesSet(); - Assert.notNull(this.configuration, "'configuration' can't be null"); - this.twitter = this.configuration.getTwitter(); - Assert.notNull(this.twitter, "'twitter' instance can't be null"); - } - @Override protected void doStop() { } diff --git a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/AbstractInboundTwitterStatusEndpointSupport.java b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/AbstractInboundTwitterStatusEndpointSupport.java index 352a16317f..7ff9e4393b 100644 --- a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/AbstractInboundTwitterStatusEndpointSupport.java +++ b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/AbstractInboundTwitterStatusEndpointSupport.java @@ -19,6 +19,7 @@ package org.springframework.integration.twitter; //import twitter4j.Status; +import org.springframework.core.task.TaskExecutor; import org.springframework.integration.twitter.model.Status; import org.springframework.integration.twitter.model.Twitter4jStatusImpl; @@ -26,6 +27,9 @@ import java.util.ArrayList; import java.util.Collections; import java.util.Comparator; import java.util.List; +import java.util.concurrent.Executor; +import java.util.concurrent.Executors; +import java.util.concurrent.ThreadFactory; /** @@ -34,8 +38,9 @@ import java.util.List; * * @author Josh Long */ -abstract public class AbstractInboundTwitterStatusEndpointSupport - extends AbstractInboundTwitterEndpointSupport { +abstract public class AbstractInboundTwitterStatusEndpointSupport extends AbstractInboundTwitterEndpointSupport { + + private Comparator statusComparator = new Comparator() { public int compare(Status status, Status status1) { diff --git a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/InboundDirectMessageStatusEndpoint.java b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/InboundDirectMessageStatusEndpoint.java index 04277ce880..3cdc2786bd 100644 --- a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/InboundDirectMessageStatusEndpoint.java +++ b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/InboundDirectMessageStatusEndpoint.java @@ -27,6 +27,9 @@ import java.util.ArrayList; import java.util.Collections; import java.util.Comparator; import java.util.List; +import java.util.concurrent.Executor; +import java.util.concurrent.Executors; +import java.util.concurrent.ThreadFactory; /** @@ -60,7 +63,13 @@ public class InboundDirectMessageStatusEndpoint extends AbstractInboundTwitterEn } @Override - protected void refresh() throws Exception { + public String getComponentType() { + return null; + } + + @Override + protected void beginPolling() throws Exception { + this.runAsAPIRateLimitsPermit(new ApiCallback() { public void run(InboundDirectMessageStatusEndpoint t, Twitter twitter) throws Exception { @@ -77,5 +86,4 @@ public class InboundDirectMessageStatusEndpoint extends AbstractInboundTwitterEn }); } - } diff --git a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/InboundMentionStatusEndpoint.java b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/InboundMentionStatusEndpoint.java index c634633ae3..83883fdf48 100644 --- a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/InboundMentionStatusEndpoint.java +++ b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/InboundMentionStatusEndpoint.java @@ -26,22 +26,26 @@ import java.util.List; * * @author Josh Long */ -public class InboundMentionStatusEndpoint - extends AbstractInboundTwitterStatusEndpointSupport { - +public class InboundMentionStatusEndpoint extends AbstractInboundTwitterStatusEndpointSupport { @Override - protected void refresh() throws Exception { - this.runAsAPIRateLimitsPermit(new ApiCallback() { - public void run(InboundMentionStatusEndpoint ctx, - Twitter twitter) throws Exception { + public String getComponentType() { + return null; + } + + @Override + protected void beginPolling() throws Exception { + this.runAsAPIRateLimitsPermit(new ApiCallback() { + + public void run(InboundMentionStatusEndpoint ctx, Twitter twitter) throws Exception { List stats = (!hasMarkedStatus()) ? twitter.getMentions() : twitter.getMentions(new Paging(ctx.getMarkerId())); - - forwardAll( fromTwitter4jStatuses( stats)); } + }); } + + } diff --git a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/InboundUpdatedStatusEndpoint.java b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/InboundUpdatedStatusEndpoint.java index 946729da18..6e3cdb6ab5 100644 --- a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/InboundUpdatedStatusEndpoint.java +++ b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/InboundUpdatedStatusEndpoint.java @@ -27,8 +27,14 @@ import twitter4j.Twitter; * @since 2.0 */ public class InboundUpdatedStatusEndpoint extends AbstractInboundTwitterStatusEndpointSupport { + @Override - protected void refresh() throws Exception { + public String getComponentType() { + return null; + } + + @Override + protected void beginPolling() throws Exception { this.runAsAPIRateLimitsPermit(new ApiCallback() { public void run(InboundUpdatedStatusEndpoint t, Twitter twitter) throws Exception { @@ -36,6 +42,5 @@ public class InboundUpdatedStatusEndpoint extends AbstractInboundTwitterStatusEn ? twitter.getFriendsTimeline() : twitter.getFriendsTimeline(new Paging(t.getMarkerId())))); } - }); - } + }); } } diff --git a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/config/DirectMessageInboundEndpointParser.java b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/config/DirectMessageInboundEndpointParser.java index 5cf79fc256..e8689282d7 100644 --- a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/config/DirectMessageInboundEndpointParser.java +++ b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/config/DirectMessageInboundEndpointParser.java @@ -21,7 +21,7 @@ public class DirectMessageInboundEndpointParser extends AbstractSingleBeanDefini @Override protected void doParse(Element element, ParserContext parserContext, BeanDefinitionBuilder builder) { - IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "channel", "requestChannel"); + IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "channel", "outputChannel"); IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "twitter-connection", "configuration"); } } diff --git a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/config/MentionInboundEndpointParser.java b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/config/MentionInboundEndpointParser.java index a0f046a09c..dec4eef6e0 100644 --- a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/config/MentionInboundEndpointParser.java +++ b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/config/MentionInboundEndpointParser.java @@ -21,7 +21,7 @@ public class MentionInboundEndpointParser extends AbstractSingleBeanDefinitionPa @Override protected void doParse(Element element, ParserContext parserContext, BeanDefinitionBuilder builder) { - IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "channel", "requestChannel"); + IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "channel", "outputChannel"); IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "twitter-connection", "configuration"); } } diff --git a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/config/UpdatedStatusInboundEndpointParser.java b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/config/UpdatedStatusInboundEndpointParser.java index 8d6d9345db..644c3e347b 100644 --- a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/config/UpdatedStatusInboundEndpointParser.java +++ b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/config/UpdatedStatusInboundEndpointParser.java @@ -22,7 +22,7 @@ public class UpdatedStatusInboundEndpointParser extends AbstractSingleBeanDefini @Override protected void doParse(Element element, ParserContext parserContext, BeanDefinitionBuilder builder) { IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, - "channel", "requestChannel"); + "channel", "outputChannel"); IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "twitter-connection", "configuration"); } diff --git a/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/TestRecievingUsingNamespace-context.xml b/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/TestRecievingUsingNamespace-context.xml index 98fdeb1d28..7e5f7ddc68 100644 --- a/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/TestRecievingUsingNamespace-context.xml +++ b/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/TestRecievingUsingNamespace-context.xml @@ -17,51 +17,50 @@ --> - - - - - + - + + + + + - + - - - + - + +